diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 78f5ab0bf..df003f4cf 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -282,6 +282,132 @@ jobs: log/ build/*/test_results/ + # The fault manager's PostgreSQL backend (-DPOSTGRES_SUPPORT=ON) on every distro. Every other job + # builds with the option OFF, its default. 14 is the oldest supported server. + postgres: + name: postgres (${{ matrix.ros_distro }}) + runs-on: ubuntu-latest + strategy: + fail-fast: false + matrix: + include: + - ros_distro: humble + os_image: ros:humble-ros-base + ccache_prefix: ccache-humble- + - ros_distro: jazzy + os_image: ros:jazzy-ros-base + # Not ccache-jazzy-: see the note in graph-watchdog. + ccache_prefix: ccache-jazzy-test- + - ros_distro: lyrical + os_image: ros:lyrical-ros-base + ccache_prefix: ccache-lyrical- + container: + image: ${{ matrix.os_image }} + # See graph-watchdog for the /dev/shm size. + options: --shm-size=1g + services: + postgres: + image: postgres:14 + env: + POSTGRES_USER: medkit + POSTGRES_PASSWORD: medkit + POSTGRES_DB: faults + options: >- + --health-cmd "pg_isready -U medkit -d faults" + --health-interval 5s + --health-timeout 5s + --health-retries 10 + timeout-minutes: 60 + defaults: + run: + shell: bash + + steps: + - name: Checkout repository + uses: actions/checkout@v4 + + - name: Install ccache + run: | + apt-get update + apt-get install -y ccache + + - name: Restore ccache + # Restore, never save, as in graph-watchdog: the distro's own sweep holds the objects of + # the dependencies, and this job adds only the fault manager and libpqxx. + uses: actions/cache/restore@v4 + with: + path: /root/.cache/ccache + key: ${{ matrix.ccache_prefix }}${{ github.sha }} + restore-keys: | + ${{ matrix.ccache_prefix }} + + - name: Install dependencies + run: | + apt-get update + # Upgrade the image to the apt sync rosdep installs from: packages from two syncs can be ABI-incompatible. + apt-get upgrade -y + # libpq-dev and git: POSTGRES_SUPPORT needs libpq and fetches libpqxx at configure time. + apt-get install -y ros-${{ matrix.ros_distro }}-test-msgs libpq-dev git + if [ "${{ matrix.ros_distro }}" = "humble" ] || [ "${{ matrix.ros_distro }}" = "jazzy" ]; then + apt-get install -y ros-${{ matrix.ros_distro }}-rmw-cyclonedds-cpp + fi + source /opt/ros/${{ matrix.ros_distro }}/setup.bash + for attempt in 1 2 3; do rosdep update && break; [ "$attempt" = 3 ] && exit 1; echo "rosdep update attempt $attempt failed, retrying" >&2; sleep 5; done + rosdep install --from-paths src --ignore-src -y \ + --skip-keys "ament_cmake_clang_tidy ament_cmake_clang_format" + + - name: Build packages + env: + CCACHE_DIR: /root/.cache/ccache + CCACHE_MAXSIZE: 500M + CCACHE_SLOPPINESS: pch_defines,time_macros + run: | + source /opt/ros/${{ matrix.ros_distro }}/setup.bash + ccache -z + colcon build --symlink-install \ + --packages-up-to ros2_medkit_fault_manager \ + --cmake-args -DCMAKE_BUILD_TYPE=Release -DPOSTGRES_SUPPORT=ON \ + --event-handlers console_direct+ + ccache -s + ./scripts/ccache_report.sh "${{ matrix.ros_distro }}-postgres" + + - name: Record the test step start + run: echo "TEST_STEP_START=$(date +%s)" >> "$GITHUB_ENV" + + - name: Run fault manager tests against PostgreSQL + timeout-minutes: 30 + env: + # The DDS choice belongs to the distro, as in graph-watchdog. + RMW_IMPLEMENTATION: ${{ (matrix.ros_distro == 'humble' || matrix.ros_distro == 'jazzy') && 'rmw_cyclonedds_cpp' || '' }} + # The job runs in a container, so the service name is the host name. + ROS2_MEDKIT_TEST_PG_CONN: postgresql://medkit:medkit@postgres:5432/faults + run: | + source /opt/ros/${{ matrix.ros_distro }}/setup.bash + colcon test --return-code-on-test-failure \ + --packages-select ros2_medkit_fault_manager \ + --ctest-args -LE linter \ + --event-handlers console_direct+ + + - name: Check the run completed, and record the environment + if: always() + uses: ./.github/actions/test-margin-and-environment + with: + start: ${{ env.TEST_STEP_START }} + cap-minutes: '30' + + - name: Show test results + if: always() + run: colcon test-result --verbose + + - name: Upload test results + if: always() + uses: actions/upload-artifact@v4 + with: + name: test-results-postgres-${{ matrix.ros_distro }} + path: | + log/ + build/*/test_results/ + # Builds AND tests Jazzy. These were two jobs (jazzy-build -> jazzy-test) # passing a tarred build/ + install/ tree between them. That split dates from # when lint ran off the same artifact in parallel; lint has since moved to diff --git a/docker/postgres-compose.yaml b/docker/postgres-compose.yaml new file mode 100644 index 000000000..091c98d40 --- /dev/null +++ b/docker/postgres-compose.yaml @@ -0,0 +1,11 @@ +services: + db: + image: postgres:18 + restart: unless-stopped + container_name: ros2_medkit_postgres_test + network_mode: host + shm_size: 128mb + environment: + POSTGRES_USER: user + POSTGRES_PASSWORD: password + POSTGRES_DB: ros2_medkit_faults_database diff --git a/docker/postgres14-compose.yaml b/docker/postgres14-compose.yaml new file mode 100644 index 000000000..952478baa --- /dev/null +++ b/docker/postgres14-compose.yaml @@ -0,0 +1,11 @@ +services: + db: + image: postgres:14 + restart: unless-stopped + container_name: ros2_medkit_postgres_test + network_mode: host + shm_size: 128mb + environment: + POSTGRES_USER: user + POSTGRES_PASSWORD: password + POSTGRES_DB: ros2_medkit_faults_database diff --git a/docs/config/fault-manager.rst b/docs/config/fault-manager.rst index 400de56c0..e0acc7f7d 100644 --- a/docs/config/fault-manager.rst +++ b/docs/config/fault-manager.rst @@ -18,8 +18,9 @@ Storage fault_manager: ros__parameters: - storage_type: "sqlite" # Storage backend: "sqlite" or "memory" + storage_type: "sqlite" # Storage backend: "sqlite", "memory" or "postgres" database_path: "/var/lib/ros2_medkit/faults.db" # Path for sqlite storage + database_url: "" # PostgreSQL connection string; empty = libpq environment variables .. list-table:: :header-rows: 1 @@ -30,10 +31,92 @@ Storage - Description * - ``storage_type`` - ``sqlite`` - - Storage backend. ``sqlite`` persists faults to disk, ``memory`` keeps in RAM only. + - Storage backend. ``sqlite`` persists faults to disk, ``memory`` keeps in RAM only, ``postgres`` + stores them in a PostgreSQL server (see `PostgreSQL Storage`_). * - ``database_path`` - ``/var/lib/ros2_medkit/faults.db`` - File path for SQLite database. Directory must exist and be writable. + * - ``database_url`` + - ``""`` + - PostgreSQL connection string, as a URI (``postgresql://db-host:5432/faults``) or as + ``key=value`` pairs. Empty takes the connection from the libpq environment variables + (``PGHOST``, ``PGPORT``, ``PGDATABASE``, ``PGUSER``, ``PGPASSWORD``) and ``~/.pgpass``. + Keep the password out of this parameter: any process on the ROS graph can read parameters. + +PostgreSQL Storage +~~~~~~~~~~~~~~~~~~ + +PostgreSQL support is off by default. To build it, install ``libpq-dev`` and configure the +package with ``-DPOSTGRES_SUPPORT=ON``. The build downloads libpqxx 7.10.7 from GitHub and links +it statically, so it needs network access at configure time. rosdep installs nothing for it. + +.. code-block:: bash + + colcon build --packages-select ros2_medkit_fault_manager --cmake-args -DPOSTGRES_SUPPORT=ON + +A build without the option stops at startup when ``storage_type`` is ``postgres``. + +- Supported servers: PostgreSQL 14 and newer. +- One database belongs to exactly one fault manager. Fault codes are the primary key, so two fault + managers on one database merge their faults and delete each other's rosbag rows. Give each fault + manager its own database, or its own schema with ``options=-csearch_path=`` in the + connection string. +- The audit log stays a local SQLite file, see ``audit_log.database_path``. +- The node never logs ``database_url``. It logs host, port, database and user, or the service + name. Error text from the server is logged and returned with the password from ``database_url`` + or ``PGPASSWORD`` replaced by ``***``. + +A wrong configuration stops the node at startup: it logs the reason and exits with code 1. A +server that cannot be reached does not: + +.. list-table:: + :header-rows: 1 + :widths: 45 55 + + * - Situation + - Fault manager + * - ``database_url`` is not a valid connection string + - Does not start. The error does not quote the string, because it can hold the password. + * - The server accepts the connection, but the schema cannot be created (for example, the + user has no ``CREATE`` privilege), or an existing table lacks a column the node uses + - Does not start. The error is the message of the server. + * - The server cannot be reached at startup: refused, timeout, unknown host, wrong password, + missing database or role + - Starts without storage and logs the reason. Requests try to connect, and the schema is + created on the first connection that succeeds. + * - The server goes away while the node runs + - Keeps running. A request that waits for the server fails after ``tcp_user_timeout``, then + requests try to reconnect. + +A wrong password or a missing database counts as unreachable. libpq reports these only as text, +and some of them pass on their own, for example a database that a container creates while it +starts. Values that libpq checks only when it connects, such as ``sslmode=required`` or +``port=abc``, also count as unreachable: the node runs without storage and logs libpq's reason on +every attempt. If the server can be reached only later and then refuses the schema, the node stops +at that point. + +When a connection is lost, the node tries twice, 500 ms apart. After a failed round, requests fail +at once for 5 seconds, then the next request tries once more. An attempt waits at most +``connect_timeout`` seconds (default 2) for each host address, so a server that does not answer +holds the node for about 2 seconds every 7 seconds. A host list, or a host name with several +addresses, multiplies that wait. Host name lookup is not part of it; give the address in +``hostaddr`` to skip it. + +On an open connection, libpq drops the connection when sent data stays unacknowledged for +``tcp_user_timeout`` milliseconds (default 5000), and TCP keepalive probes start after +``keepalives_idle`` seconds of silence (default 5). This bounds a request to a server that went +off the network. A server process that stops answering while its host still acknowledges packets +is not bounded: the request, and with it the node, waits for that server. + +Set these values in ``database_url``; ``PGCONNECT_TIMEOUT`` also sets ``connect_timeout``. With a +libpq service (``service=`` or ``PGSERVICE``) the node adds none of these defaults, so set them in +the service file. + +While there is no storage, a service that has an error field answers ``Fault storage unavailable``. +``ListFaults`` has no error field and answers with an empty list. The near-miss trim and the +reclassification of HEALED faults run only at startup, so they are skipped when the server cannot +be reached then. Snapshot capture starts with the first answer of the server. The bags of a fault +cleared while the server cannot be reached are not deleted then. Debounce Settings ~~~~~~~~~~~~~~~~~ @@ -581,7 +664,9 @@ by default: with it off there is no table, no file and no write cost. * - ``audit_log.database_path`` - ``""`` - Where the audit database lives. Empty puts it beside the fault database, - or in memory when the fault store is itself in memory or not SQLite. + or in memory when the fault store is in memory or of an unknown type. With + ``storage_type: postgres`` it is a local SQLite file ``fault_audit.db`` next to + ``database_path``, so the audit trail stays on the robot. Correlation Configuration ------------------------- @@ -624,6 +709,7 @@ Complete Example # Storage storage_type: "sqlite" database_path: "/var/lib/ros2_medkit/faults.db" + database_url: "" # Debounce for a reporter that repeats its events while a condition holds: # three FAILED events confirm, and four PASSED events heal from there. diff --git a/src/ros2_medkit_fault_manager/CMakeLists.txt b/src/ros2_medkit_fault_manager/CMakeLists.txt index 064cf44f1..6a2ef6c07 100644 --- a/src/ros2_medkit_fault_manager/CMakeLists.txt +++ b/src/ros2_medkit_fault_manager/CMakeLists.txt @@ -12,9 +12,11 @@ # See the License for the specific language governing permissions and # limitations under the License. -cmake_minimum_required(VERSION 3.8) +cmake_minimum_required(VERSION 3.14) project(ros2_medkit_fault_manager) +option(POSTGRES_SUPPORT "Enable PostgreSQL support" OFF) + set(CMAKE_CXX_STANDARD 17) set(CMAKE_CXX_STANDARD_REQUIRED ON) set(CMAKE_EXPORT_COMPILE_COMMANDS ON) @@ -32,6 +34,75 @@ find_package(ament_cmake REQUIRED) find_package(rclcpp REQUIRED) find_package(ros2_medkit_msgs REQUIRED) find_package(ros2_medkit_serialization REQUIRED) +if(POSTGRES_SUPPORT) + find_package(PostgreSQL REQUIRED) + include(FetchContent) + + # EXCLUDE_FROM_ALL keeps libpqxx out of this package's install. CMake < 3.28 has no such + # option in fetchcontent_declare, so the subdirectory is added by hand below. + set(_pqxx_exclude_from_all "") + if(CMAKE_VERSION VERSION_GREATER_EQUAL 3.28) + set(_pqxx_exclude_from_all EXCLUDE_FROM_ALL) + endif() + fetchcontent_declare(pqxx + GIT_REPOSITORY https://github.com/jtv/libpqxx.git + GIT_TAG 7.10.7 + GIT_SHALLOW TRUE + ${_pqxx_exclude_from_all} + ) + + # Save medkit flags to restore them later. Not required per-se, but doesn't hurt to handle them either + set(_medkit_saved_cxx_flags "${CMAKE_CXX_FLAGS}") + get_directory_property(_medkit_saved_opts COMPILE_OPTIONS) + string(REGEX REPLACE "(^|[ ])-W[^ ]*" " " CMAKE_CXX_FLAGS "${CMAKE_CXX_FLAGS}") + set_directory_properties(PROPERTIES COMPILE_OPTIONS "") + + # No warnings-as-errors, include-what-you-use or clang-tidy for pqxx. An empty normal variable + # hides a cache value from -D; unset() would expose it. Restored after pqxx. + set(_medkit_saved_warning_as_error "${CMAKE_COMPILE_WARNING_AS_ERROR}") + set(CMAKE_COMPILE_WARNING_AS_ERROR OFF) + set(CMAKE_CXX_INCLUDE_WHAT_YOU_USE "") + set(CMAKE_CXX_CLANG_TIDY "") + # Keep the libpqxx tests out of this build. Its BUILD_DOC is OFF by default; setting it here makes + # CMake 3.22 warn about CMP0077. + set(SKIP_BUILD_TEST ON) + + # libpqxx must be static: it is not installed. ament on Lyrical defaults BUILD_SHARED_LIBS to ON. + set(_medkit_saved_build_shared_libs "${BUILD_SHARED_LIBS}") + set(BUILD_SHARED_LIBS OFF) + if(CMAKE_VERSION VERSION_GREATER_EQUAL 3.28) + fetchcontent_makeavailable(pqxx) + else() + fetchcontent_getproperties(pqxx) + if(NOT pqxx_POPULATED) + fetchcontent_populate(pqxx) + add_subdirectory(${pqxx_SOURCE_DIR} ${pqxx_BINARY_DIR} EXCLUDE_FROM_ALL) + endif() + endif() + set(BUILD_SHARED_LIBS "${_medkit_saved_build_shared_libs}") + set(CMAKE_COMPILE_WARNING_AS_ERROR "${_medkit_saved_warning_as_error}") + unset(CMAKE_CXX_INCLUDE_WHAT_YOU_USE) + unset(CMAKE_CXX_CLANG_TIDY) + + # Restore medkit flags now that pqxx is compiled + set(CMAKE_CXX_FLAGS "${_medkit_saved_cxx_flags}") + set_directory_properties(PROPERTIES COMPILE_OPTIONS "${_medkit_saved_opts}") + + # Make headers available to SYSTEM + # Once CMake3.25 is globally supported (Ubuntu 22.04 ships with CMake3.22) + # the snippet below could be replaced with fetchcontent_declare(pqxx ... SYSTEM) + foreach(_pqxx_tgt pqxx pqxx_shared pqxx_static) + if(TARGET ${_pqxx_tgt}) + target_compile_options(${_pqxx_tgt} PRIVATE -w) + get_target_property(_pqxx_inc ${_pqxx_tgt} INTERFACE_INCLUDE_DIRECTORIES) + if(_pqxx_inc) + set_target_properties(${_pqxx_tgt} PROPERTIES + INTERFACE_SYSTEM_INCLUDE_DIRECTORIES "${_pqxx_inc}") + endif() + endif() + endforeach() +endif() + find_package(SQLite3 REQUIRED) find_package(nlohmann_json REQUIRED) # OpenSSL EVP SHA-256 for the tamper-evident audit log hash chain @@ -45,19 +116,27 @@ find_package(rosbag2_storage REQUIRED) medkit_detect_compat_defs() # Library target (shared between executable and tests) +set(FAULT_MANAGER_FILES + src/capture_thread_pool.cpp + src/fault_manager_node.cpp + src/fault_storage.cpp + src/sqlite_fault_storage.cpp + src/fault_audit_log.cpp + src/snapshot_capture.cpp + src/rosbag_capture.cpp + src/correlation/types.cpp + src/correlation/config_parser.cpp + src/correlation/pattern_matcher.cpp + src/correlation/correlation_engine.cpp + src/entity_threshold_resolver.cpp +) +if(POSTGRES_SUPPORT) + list(APPEND FAULT_MANAGER_FILES + src/postgres_fault_storage.cpp + ) +endif() add_library(fault_manager_lib STATIC - src/capture_thread_pool.cpp - src/fault_manager_node.cpp - src/fault_storage.cpp - src/sqlite_fault_storage.cpp - src/fault_audit_log.cpp - src/snapshot_capture.cpp - src/rosbag_capture.cpp - src/correlation/types.cpp - src/correlation/config_parser.cpp - src/correlation/pattern_matcher.cpp - src/correlation/correlation_engine.cpp - src/entity_threshold_resolver.cpp + ${FAULT_MANAGER_FILES} ) target_include_directories(fault_manager_lib PUBLIC @@ -73,11 +152,22 @@ medkit_target_dependencies(fault_manager_lib PUBLIC rosbag2_storage ) +set(FAULT_MANAGER_TARGET_LIBS + SQLite::SQLite3 + nlohmann_json::nlohmann_json + yaml-cpp::yaml-cpp + OpenSSL::Crypto +) + +if(POSTGRES_SUPPORT) + # libpq directly as well: the connection string is checked with PQconninfoParse. + list(APPEND FAULT_MANAGER_TARGET_LIBS + pqxx + PostgreSQL::PostgreSQL + ) +endif() target_link_libraries(fault_manager_lib PUBLIC - SQLite::SQLite3 - nlohmann_json::nlohmann_json - yaml-cpp::yaml-cpp - OpenSSL::Crypto + ${FAULT_MANAGER_TARGET_LIBS} ) medkit_apply_compat_defs(fault_manager_lib) @@ -136,6 +226,13 @@ if(BUILD_TESTING) target_link_libraries(test_sqlite_storage fault_manager_lib) medkit_target_dependencies(test_sqlite_storage rclcpp ros2_medkit_msgs) + # PostgreSQL storage tests +if(POSTGRES_SUPPORT) + medkit_add_gtest(test_postgres_storage test/test_postgres_storage.cpp) + target_link_libraries(test_postgres_storage fault_manager_lib) + medkit_target_dependencies(test_postgres_storage rclcpp ros2_medkit_msgs) +endif() + # Rosbag retention parity: every assertion runs against both storage backends. medkit_add_gtest(test_rosbag_storage_parity test/test_rosbag_storage_parity.cpp) target_link_libraries(test_rosbag_storage_parity fault_manager_lib) @@ -232,4 +329,8 @@ if(BUILD_TESTING) ros2_medkit_relax_vendor_warnings() endif() +# Defined only when ON: the sources test it with #ifdef. +if(POSTGRES_SUPPORT) + target_compile_definitions(fault_manager_lib PUBLIC POSTGRES_SUPPORT) +endif() ament_package() diff --git a/src/ros2_medkit_fault_manager/README.md b/src/ros2_medkit_fault_manager/README.md index 977f5a8c4..e6ce3af93 100644 --- a/src/ros2_medkit_fault_manager/README.md +++ b/src/ros2_medkit_fault_manager/README.md @@ -37,11 +37,11 @@ ros2 service call /fault_manager/clear_fault ros2_medkit_msgs/srv/ClearFault \ ## Services -| Service | Type | Description | -|---------|------|-------------| -| `~/report_fault` | `ros2_medkit_msgs/srv/ReportFault` | Report a fault occurrence | -| `~/list_faults` | `ros2_medkit_msgs/srv/ListFaults` | Query faults with filtering | -| `~/clear_fault` | `ros2_medkit_msgs/srv/ClearFault` | Clear/acknowledge a fault | +| Service | Type | Description | +| ----------------- | ----------------------------------- | ------------------------------- | +| `~/report_fault` | `ros2_medkit_msgs/srv/ReportFault` | Report a fault occurrence | +| `~/list_faults` | `ros2_medkit_msgs/srv/ListFaults` | Query faults with filtering | +| `~/clear_fault` | `ros2_medkit_msgs/srv/ClearFault` | Clear/acknowledge a fault | | `~/get_snapshots` | `ros2_medkit_msgs/srv/GetSnapshots` | Get topic snapshots for a fault | ## Features @@ -50,7 +50,7 @@ ros2 service call /fault_manager/clear_fault ros2_medkit_msgs/srv/ClearFault \ - **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 - **Severity escalation**: Fault severity is updated if a higher severity is reported -- **Persistent storage**: SQLite backend ensures faults survive node restarts +- **Persistent storage**: SQLite (default) or PostgreSQL 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 @@ -60,16 +60,17 @@ ros2 service call /fault_manager/clear_fault ros2_medkit_msgs/srv/ClearFault \ ## Parameters -| Parameter | Type | Default | Description | -|-----------|------|---------|-------------| -| `storage_type` | string | `"sqlite"` | Storage backend: `"sqlite"` or `"memory"` | -| `database_path` | string | `"/var/lib/ros2_medkit/faults.db"` | Path to SQLite database file | -| `confirmation_threshold` | int | `-1` | Counter value at which faults are confirmed | -| `healing_enabled` | bool | `false` | Enable automatic healing via PASSED events | -| `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) | +| Parameter | Type | Default | Description | +| ------------------------------- | ------ | ---------------------------------- | ------------------------------------------------------------------------------- | +| `storage_type` | string | `"sqlite"` | Storage backend: `"sqlite"`, `"memory"` or `"postgres"` | +| `database_path` | string | `"/var/lib/ros2_medkit/faults.db"` | Path to SQLite database file | +| `database_url` | string | `""` | PostgreSQL connection string; empty = libpq environment variables | +| `confirmation_threshold` | int | `-1` | Counter value at which faults are confirmed | +| `healing_enabled` | bool | `false` | Enable automatic healing via PASSED events | +| `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) | ### Snapshot Parameters @@ -79,28 +80,30 @@ Each confirm also writes a **freeze-frame**: a single compact JSON object mappin 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. -That single-writer property also shapes what each fault of a burst gets. Nothing is buffered while a post-fault window is open (messages go straight into the open bag), and the flush that opened that bag already emptied the ring buffer, so a fault confirming right *after* the window closes has no pre-fault history available. What a confirmation gets is decided by the buffer, so a fault arriving before any captured topic has published lands the same way. It gets a **post-fault-only bag**: its own recording holding just its `duration_after_sec` window, entered through the same post-roll state machine, so later faults of the burst attach to it normally. With `duration_after_sec: 0` there is no window to record into and such a fault gets no bag; if the bag cannot be written at all, no recording is opened and no metadata row is stored. The `duration_sec` on a stored bag is the span the recording was open rather than the configured windows, so a post-fault-only bag usually reports roughly `duration_after_sec` where a full one reports its buffered history too. It is a recording span, not a content span: a window during which nothing was published still reports the seconds it covered. It can also exceed `duration_sec + duration_after_sec`, because the ring buffer is pruned only when a message arrives - a topic that stops publishing keeps its last window buffered until the next confirmation flushes it, which is deliberate for a black box. See `docs/config/fault-manager.rst` for the full lifecycle. - -| Parameter | Type | Default | Description | -|-----------|------|---------|-------------| -| `snapshots.enabled` | bool | `true` | Enable/disable snapshot capture | -| `snapshots.background_capture` | bool | `false` | Use background subscriptions (caches latest message) vs on-demand capture | -| `snapshots.timeout_sec` | double | `1.0` | Timeout waiting for topic message (on-demand mode) | -| `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.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). | -| `snapshots.capture_queue_full_policy` | string | `reject_newest` | Policy when the queue is full: `reject_newest` or `drop_oldest`. | +That single-writer property also shapes what each fault of a burst gets. Nothing is buffered while a post-fault window is open (messages go straight into the open bag), and the flush that opened that bag already emptied the ring buffer, so a fault confirming right _after_ the window closes has no pre-fault history available. What a confirmation gets is decided by the buffer, so a fault arriving before any captured topic has published lands the same way. It gets a **post-fault-only bag**: its own recording holding just its `duration_after_sec` window, entered through the same post-roll state machine, so later faults of the burst attach to it normally. With `duration_after_sec: 0` there is no window to record into and such a fault gets no bag; if the bag cannot be written at all, no recording is opened and no metadata row is stored. The `duration_sec` on a stored bag is the span the recording was open rather than the configured windows, so a post-fault-only bag usually reports roughly `duration_after_sec` where a full one reports its buffered history too. It is a recording span, not a content span: a window during which nothing was published still reports the seconds it covered. It can also exceed `duration_sec + duration_after_sec`, because the ring buffer is pruned only when a message arrives - a topic that stops publishing keeps its last window buffered until the next confirmation flushes it, which is deliberate for a black box. See `docs/config/fault-manager.rst` for the full lifecycle. + +| Parameter | Type | Default | Description | +| ------------------------------------- | -------- | --------------- | -------------------------------------------------------------------------------------------------------------------------- | +| `snapshots.enabled` | bool | `true` | Enable/disable snapshot capture | +| `snapshots.background_capture` | bool | `false` | Use background subscriptions (caches latest message) vs on-demand capture | +| `snapshots.timeout_sec` | double | `1.0` | Timeout waiting for topic message (on-demand mode) | +| `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.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). | +| `snapshots.capture_queue_full_policy` | string | `reject_newest` | Policy when the queue is full: `reject_newest` or `drop_oldest`. | **Topic Resolution Priority:** + 1. `fault_specific` - Exact match for fault code (configured via YAML config file) 2. `patterns` - Regex pattern match (configured via YAML config file) 3. `default_topics` - Fallback for all faults **Example YAML config file** (`snapshots.yaml`): + ```yaml fault_specific: MOTOR_OVERHEAT: @@ -135,6 +138,18 @@ format used by black-box capture (`snapshots.rosbag.format`, see Rosbag Capture **Memory**: Faults are stored in memory only. Useful for testing or when persistence is not required. +**PostgreSQL**: Faults are stored in an external PostgreSQL server (14 or newer) and survive node restarts. The audit log stays a local SQLite file next to `database_path`. One database belongs to exactly one fault manager; give each fault manager its own database or schema. + +PostgreSQL support is off by default. To build it, install `libpq-dev` and pass `-DPOSTGRES_SUPPORT=ON`: + +```bash +colcon build --packages-select ros2_medkit_fault_manager --cmake-args -DPOSTGRES_SUPPORT=ON +``` + +The build downloads libpqxx 7.10.7 from GitHub at configure time and links it statically; rosdep installs nothing for it. Keep the password out of `database_url`: leave it empty and set `PGHOST`, `PGPORT`, `PGDATABASE`, `PGUSER` and `PGPASSWORD` (or use `~/.pgpass`). The node never logs the connection string, and it removes the password from error text it logs. + +A wrong configuration stops the node at startup with exit code 1: a connection string libpq cannot parse, or a server that accepts the connection but does not allow creating the schema, or a table that lacks a column the node uses. A server that cannot be reached (including a wrong password or a missing database) does not stop it: the node runs without storage, services with an error field answer `Fault storage unavailable`, and `ReportFault` answers `accepted=false`. After a failed connection attempt, requests fail at once for 5 s and the next request tries again. An attempt waits at most 2 s per host address unless `connect_timeout` or `PGCONNECT_TIMEOUT` says otherwise, and a request on an open connection to a server that went off the network fails after `tcp_user_timeout` (5 s). With a libpq service, these values come from the service file. See `docs/config/fault-manager.rst` for the full table. + ## Near-Miss Series A **near miss** is a FAILED report that moved the debounce counter without the fault ending up @@ -197,21 +212,21 @@ An optional append-only, hash-chained audit log records every fault state transi 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). -**Completeness is an integrity property.** `verify()` proves nothing was *deleted* from the chain, but it cannot prove a transition that was *never appended*. So a silently dropped append is a hole `verify()` can never see. Every transition on the write path is therefore audited (occurred, timer/threshold confirmations, auto-heal, and clears), and an append failure is never swallowed silently: it increments a dropped-writes counter and clears an "audit healthy" flag. **These are in-process signals only** (C++ getters on the node). This revision exposes no service/REST/health endpoint that surfaces audit health or lets an operator run `verify()` at runtime, so the signals are **not operator-observable at runtime yet** - a runtime read/verify/health surface is future work. With `audit_log.fail_closed` set, an append failure is re-raised as a **fail-FAST** error so a compliance-strict deployment learns the audit broke. This does **not** roll back the fault-state change that already committed: the audit log and the fault store are **separate SQLite databases**, so there is no cross-DB atomicity, and `fail_closed` is a broken-audit alarm requiring operator action, not a rollback. The default (`fail_closed=false`) keeps fault processing running; either way the in-process signals record the gap. +**Completeness is an integrity property.** `verify()` proves nothing was _deleted_ from the chain, but it cannot prove a transition that was _never appended_. So a silently dropped append is a hole `verify()` can never see. Every transition on the write path is therefore audited (occurred, timer/threshold confirmations, auto-heal, and clears), and an append failure is never swallowed silently: it increments a dropped-writes counter and clears an "audit healthy" flag. **These are in-process signals only** (C++ getters on the node). This revision exposes no service/REST/health endpoint that surfaces audit health or lets an operator run `verify()` at runtime, so the signals are **not operator-observable at runtime yet** - a runtime read/verify/health surface is future work. With `audit_log.fail_closed` set, an append failure is re-raised as a **fail-FAST** error so a compliance-strict deployment learns the audit broke. This does **not** roll back the fault-state change that already committed: the audit log and the fault store are **separate SQLite databases**, so there is no cross-DB atomicity, and `fail_closed` is a broken-audit alarm requiring operator action, not a rollback. The default (`fail_closed=false`) keeps fault processing running; either way the in-process signals record the gap. -`verify()` walks the persisted chain oldest-first and recomputes every link: editing a row breaks its `record_hash`, and deleting a row breaks the next row's `prev_hash` linkage. Deleting the newest row *while leaving the head untouched* is caught by the persisted-head check (the head is read straight from the DB). However, deleting the newest row(s) **and** repointing the head with a single `UPDATE audit_chain_head SET seq=..., record_hash=...` to the prior row's values costs no more than any other casual edit - it is **not** the "recompute the entire chain" the threat model below might suggest - and the truncated chain still verifies. There is no external record that a later `seq` ever existed, so this tail-truncation is undetectable by design. +`verify()` walks the persisted chain oldest-first and recomputes every link: editing a row breaks its `record_hash`, and deleting a row breaks the next row's `prev_hash` linkage. Deleting the newest row _while leaving the head untouched_ is caught by the persisted-head check (the head is read straight from the DB). However, deleting the newest row(s) **and** repointing the head with a single `UPDATE audit_chain_head SET seq=..., record_hash=...` to the prior row's values costs no more than any other casual edit - it is **not** the "recompute the entire chain" the threat model below might suggest - and the truncated chain still verifies. There is no external record that a later `seq` ever existed, so this tail-truncation is undetectable by design. -**Threat model (read this).** The chain is **unkeyed**, and the head and segment anchors live in the **same writable SQLite file** as the rows. `verify()` therefore catches edits or deletions that did **not** also recompute the chain - that is, casual or accidental tampering, and the bookkeeping bugs that would otherwise lose records. The append-only triggers are defense-in-depth: `audit_log` rejects out-of-band UPDATE/DELETE, `audit_anchors` carries the same guard-gated triggers so an out-of-band INSERT/UPDATE/DELETE of an anchor is rejected too, and the rotation-prune guard (`audit_prune_guard`) is itself protected by a trigger so an external writer cannot simply flip it open and then delete a prefix (or forge an anchor) - that flip is only permitted from the in-process connection that holds a per-connection temp marker. The single-row chain head (`audit_chain_head`) is intentionally **not** trigger-protected (a trigger there would block the legitimate head update inside the append transaction); a casual edit or delete of the head is instead caught by `verify()` via the seq/hash/head-mismatch checks. None of this stops an attacker with write access to the file: such an attacker can create the same temp marker or drop the triggers, and recompute the entire chain (head and anchors included) to forge a self-consistent history - and cheaper still, the tail-truncation above and the forged prefix-truncation below need no recompute at all. The triggers are **not** a security boundary - this is tamper-**evident**, not tamper-**proof**. True tamper-*proofing* requires a key or signature over the head (so it cannot be recomputed without the key) or external anchoring of the head hash to an append-only store you do not control; both are out of scope here and belong to the audit-log exporter / signing follow-up. +**Threat model (read this).** The chain is **unkeyed**, and the head and segment anchors live in the **same writable SQLite file** as the rows. `verify()` therefore catches edits or deletions that did **not** also recompute the chain - that is, casual or accidental tampering, and the bookkeeping bugs that would otherwise lose records. The append-only triggers are defense-in-depth: `audit_log` rejects out-of-band UPDATE/DELETE, `audit_anchors` carries the same guard-gated triggers so an out-of-band INSERT/UPDATE/DELETE of an anchor is rejected too, and the rotation-prune guard (`audit_prune_guard`) is itself protected by a trigger so an external writer cannot simply flip it open and then delete a prefix (or forge an anchor) - that flip is only permitted from the in-process connection that holds a per-connection temp marker. The single-row chain head (`audit_chain_head`) is intentionally **not** trigger-protected (a trigger there would block the legitimate head update inside the append transaction); a casual edit or delete of the head is instead caught by `verify()` via the seq/hash/head-mismatch checks. None of this stops an attacker with write access to the file: such an attacker can create the same temp marker or drop the triggers, and recompute the entire chain (head and anchors included) to forge a self-consistent history - and cheaper still, the tail-truncation above and the forged prefix-truncation below need no recompute at all. The triggers are **not** a security boundary - this is tamper-**evident**, not tamper-**proof**. True tamper-_proofing_ requires a key or signature over the head (so it cannot be recomputed without the key) or external anchoring of the head hash to an append-only store you do not control; both are out of scope here and belong to the audit-log exporter / signing follow-up. -**Retention/rotation**: when more than `audit_log.retention_max_records` rows are retained, the oldest segment is *sealed* (its final `seq` + hash are persisted as an anchor) and then pruned. The surviving tail still verifies because the oldest retained row links back to the sealed anchor. Only the anchor at the current prune boundary is kept - the same rotation drops older anchors - so `audit_anchors` stays bounded (one row) instead of growing one row per rotation. Because `verify()` treats any matching sealed anchor as a valid tail root, a **forged** prefix-truncation (an out-of-band actor deletes a prefix and inserts a matching anchor) is **indistinguishable** from legitimate pruning: "the surviving tail still verifies" therefore covers a forged truncation exactly as well as a real one. The guard-gated `audit_anchors` triggers raise the bar for this (casual/accidental only) but, like every trigger here, a write-capable adversary can drop them - so this stays tamper-**evident**, not tamper-**proof**. +**Retention/rotation**: when more than `audit_log.retention_max_records` rows are retained, the oldest segment is _sealed_ (its final `seq` + hash are persisted as an anchor) and then pruned. The surviving tail still verifies because the oldest retained row links back to the sealed anchor. Only the anchor at the current prune boundary is kept - the same rotation drops older anchors - so `audit_anchors` stays bounded (one row) instead of growing one row per rotation. Because `verify()` treats any matching sealed anchor as a valid tail root, a **forged** prefix-truncation (an out-of-band actor deletes a prefix and inserts a matching anchor) is **indistinguishable** from legitimate pruning: "the surviving tail still verifies" therefore covers a forged truncation exactly as well as a real one. The guard-gated `audit_anchors` triggers raise the bar for this (casual/accidental only) but, like every trigger here, a write-capable adversary can drop them - so this stays tamper-**evident**, not tamper-**proof**. -| Parameter | Type | Default | Description | -|-----------|------|---------|-------------| -| `audit_log.enabled` | bool | `false` | Enable the tamper-evident audit log | -| `audit_log.transitions` | string | `"all"` | Which transitions to record: `"all"` (occurred/confirmed/healed/cleared) or `"confirmed_only"`. Lifecycle markers are always recorded. | -| `audit_log.database_path` | string | `""` | SQLite path. Empty => sibling `fault_audit.db` next to the fault DB (or `:memory:` for in-memory fault stores) | -| `audit_log.retention_max_records` | int | `0` | Seal + prune the oldest segment beyond this many retained records (0 = unlimited) | -| `audit_log.fail_closed` | bool | `false` | When `true`, an audit append failure is re-raised as a fail-FAST error signalling the audit chain is broken and needs operator action. It does **not** roll back the already-committed fault-state change (the fault store is a separate DB - no cross-DB atomicity). When `false`, the failure is logged and counted and fault processing continues. Either way the gap is recorded via the in-process dropped-writes / audit-healthy signals (not operator-observable at runtime yet). | +| Parameter | Type | Default | Description | +| --------------------------------- | ------ | ------- | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `audit_log.enabled` | bool | `false` | Enable the tamper-evident audit log | +| `audit_log.transitions` | string | `"all"` | Which transitions to record: `"all"` (occurred/confirmed/healed/cleared) or `"confirmed_only"`. Lifecycle markers are always recorded. | +| `audit_log.database_path` | string | `""` | SQLite path. Empty => sibling `fault_audit.db` next to the fault DB (or `:memory:` for in-memory fault stores) | +| `audit_log.retention_max_records` | int | `0` | Seal + prune the oldest segment beyond this many retained records (0 = unlimited) | +| `audit_log.fail_closed` | bool | `false` | When `true`, an audit append failure is re-raised as a fail-FAST error signalling the audit chain is broken and needs operator action. It does **not** roll back the already-committed fault-state change (the fault store is a separate DB - no cross-DB atomicity). When `false`, the failure is logged and counted and fault processing continues. Either way the gap is recorded via the in-process dropped-writes / audit-healthy signals (not operator-observable at runtime yet). | ## Usage @@ -374,17 +389,17 @@ ros2 run ros2_medkit_fault_manager fault_manager_node --ros-args \ ### Correlation Parameters -| Parameter | Type | Default | Description | -|-----------|------|---------|-------------| -| `correlation.config_file` | string | `""` | Path to correlation YAML config (empty = disabled) | -| `correlation.cleanup_interval_sec` | double | `5.0` | Interval for cleaning up expired pending correlations (seconds) | +| Parameter | Type | Default | Description | +| ---------------------------------- | ------ | ------- | --------------------------------------------------------------- | +| `correlation.config_file` | string | `""` | Path to correlation YAML config (empty = disabled) | +| `correlation.cleanup_interval_sec` | double | `5.0` | Interval for cleaning up expired pending correlations (seconds) | ### Configuration File Format ```yaml correlation: enabled: true - default_window_ms: 500 # Default time window for symptom detection + default_window_ms: 500 # Default time window for symptom detection # Reusable fault patterns (supports wildcards with *) patterns: @@ -405,8 +420,8 @@ correlation: symptoms: - pattern: motor_errors - pattern: drive_faults - window_ms: 1000 # Symptoms within 1s of root cause - mute_symptoms: true # Don't publish symptom events + window_ms: 1000 # Symptoms within 1s of root cause + mute_symptoms: true # Don't publish symptom events auto_clear_with_root: true # Clear symptoms when root cause clears # Auto-cluster rule: Group communication errors @@ -415,15 +430,16 @@ correlation: mode: auto_cluster match: - pattern: comm_errors - min_count: 3 # Need 3 faults to form cluster - window_ms: 500 # Within 500ms - show_as_single: true # Only show representative fault - representative: highest_severity # first | most_recent | highest_severity + min_count: 3 # Need 3 faults to form cluster + window_ms: 500 # Within 500ms + show_as_single: true # Only show representative fault + representative: highest_severity # first | most_recent | highest_severity ``` ### Pattern Wildcards Patterns support `*` wildcard matching: + - `MOTOR_*` matches `MOTOR_COMM`, `MOTOR_TIMEOUT`, `MOTOR_DRIVE_FAULT` - `*_COMM_*` matches `MOTOR_COMM_FL`, `SENSOR_COMM_TIMEOUT` - `*_TIMEOUT` matches `MOTOR_TIMEOUT`, `SENSOR_TIMEOUT` @@ -439,6 +455,7 @@ 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`) @@ -447,10 +464,12 @@ Response includes: ### REST API (via Gateway) Query parameters for GET `/api/v1/faults`: + - `include_muted=true`: Include muted fault details in response - `include_clusters=true`: Include cluster details in response Response fields: + ```json { "faults": [...], @@ -482,6 +501,7 @@ Response fields: ``` When clearing a root cause fault, `auto_cleared_codes` lists symptoms that were auto-cleared: + ```json { "status": "success", 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..5a4806163 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 @@ -205,6 +205,7 @@ class FaultManagerNode : public rclcpp::Node { std::string storage_type_; std::string database_path_; + std::string database_url_; int32_t confirmation_threshold_{-1}; bool healing_enabled_{false}; int32_t healing_threshold_{3}; 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..645978b1b 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 @@ -15,6 +15,7 @@ #pragma once #include +#include #include #include #include @@ -198,6 +199,22 @@ struct RosbagFileInfo { /// Abstract interface for fault storage backends class FaultStorage { public: + /// Thrown by a backend that cannot reach its database server. The node keeps running and answers with an error. + class IgnorableConnectionException : public std::exception { + protected: + std::string message; + + public: + explicit IgnorableConnectionException(const std::string & msg = "IgnorableConnectionException") : message(msg) { + } + + const char * what() const noexcept override { + return message.c_str(); + } + + virtual ~IgnorableConnectionException() = default; + }; + virtual ~FaultStorage() = default; /// Set debounce configuration diff --git a/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/postgres_fault_storage.hpp b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/postgres_fault_storage.hpp new file mode 100644 index 000000000..466b5402d --- /dev/null +++ b/src/ros2_medkit_fault_manager/include/ros2_medkit_fault_manager/postgres_fault_storage.hpp @@ -0,0 +1,203 @@ +// Copyright 2026 gstavrinos +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +#pragma once + +#include + +#include +#include +#include +#include +#include +#include + +#include "ros2_medkit_fault_manager/fault_storage.hpp" + +namespace ros2_medkit_fault_manager { + +/// PostgreSQL-based fault storage implementation with persistence +/// Thread-safe implementation using mutex protection on connection access +class PgFaultStorage : public FaultStorage { + public: + using FaultStorage::IgnorableConnectionException; + /// Create PostgreSQL fault storage. An unreachable server is not an error: every call then throws + /// IgnorableConnectionException until a connection succeeds, and the schema is created then. + /// @param conn_info libpq connection string or URI; empty uses the libpq environment variables + /// @throws std::invalid_argument if libpq rejects @p conn_info + /// @throws std::runtime_error if the server accepts the connection but the schema cannot be created + explicit PgFaultStorage(const std::string & conn_info); + + /// Same as above, with the retry policy. + /// @param conn_info libpq connection string or URI; empty uses the libpq environment variables + /// @param max_retries Attempts after the first one, for a connect and for a transaction. Must be >= 0. + /// @param reconnection_delay_ms Delay in milliseconds between two attempts + /// @throws std::invalid_argument if libpq rejects @p conn_info + /// @throws std::runtime_error if the server accepts the connection but the schema cannot be created + explicit PgFaultStorage(const std::string & conn_info, const int max_retries, const unsigned reconnection_delay_ms); + + /// Destructor - closes database connection + ~PgFaultStorage() override; + + // Non-copyable, non-movable (owns PostgreSQL connection) + PgFaultStorage(const PgFaultStorage &) = delete; + PgFaultStorage & operator=(const PgFaultStorage &) = delete; + PgFaultStorage(PgFaultStorage &&) = delete; + PgFaultStorage & operator=(PgFaultStorage &&) = delete; + + void set_debounce_config(const DebounceConfig & config) override; + DebounceConfig get_debounce_config() const override; + + 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) override; + + 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; + + bool clear_fault(const std::string & fault_code) override; + + size_t size() const override; + + bool contains(const std::string & fault_code) const 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; + bool retains_snapshots_on_clear() const override; + + void set_max_rosbags_per_fault(size_t max_count) override; + + 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; + 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; + 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; + + 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::vector get_rosbag_files_by_recording(const std::string & recording_id) const override; + bool delete_rosbag_file(const std::string & fault_code) 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 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; + + /// Get the connection info string used to initialize the database + const std::string & conn_info() const { + return conn_info_; + } + + /// "host=... port=... dbname=... user=..." from database_url and the libpq defaults. With a service, + /// "service=..." and the keys database_url sets. Never the password. + std::string target() const; + + /// Whether a connection to the server is open. + bool connected() const; + + private: + /// Wrapper the pqxx exec function + template + pqxx::result execute(pqxx::work & tx, const std::string & query, Args &&... args) const; + + /// Open a new connection unless the current one is open. Caller holds mutex_. + void ensure_connection() const; + + /// Run @p fn in one transaction and commit. On a broken connection, reconnect and run it again. Caller holds mutex_. + template + auto run_in_transaction(const char * what, Fn && fn) const -> decltype(fn(std::declval())); + + /// Create the tables and indexes that do not exist yet, in @p tx, and check the columns of existing ones. + void create_schema(pqxx::work & tx) const; + + /// @p text with every copy of the known password replaced by "***". + std::string redact(std::string text) const; + + /// Drop the connection and start the backoff after a failed round. Caller holds mutex_. + void fail_round(const std::string & error) const; + + /// Whether any fault at all still references @p file_path. Caller holds mutex_. + bool path_referenced(const std::string & file_path) const; + + /// store_rosbag_file body without taking mutex_. Caller holds mutex_ and + /// manages transaction scope. Returns replaced bag path if applicable. + std::vector store_rosbag_file_locked(const RosbagFileInfo & info, pqxx::work & tx); + + /// report_fault_event body without taking mutex_ or opening a transaction. Caller holds mutex_ + /// and supplies the transaction, so the fault row and any near-miss row commit together. + bool report_fault_event_locked(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, pqxx::work & tx); + + /// Append one entry to the near-miss series and evict the oldest entries beyond + /// max_near_misses_per_fault_. Caller holds mutex_ and has already written the fault row. + /// @param fault_code The fault code that nearly confirmed + /// @param occurred_at_ns Timestamp of the report + /// @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 resulting_status Fault status after the report was applied + /// @param tx Transaction the fault row was written in + 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, + const std::string & resulting_status, pqxx::work & tx); + + /// Deserialize JSON array string from PostgreSQL TEXT/JSONB field + static std::vector parse_json_array(const std::string & json_str); + + /// Serialize vector of strings to JSON array string (for PostgreSQL JSONB) + static std::string serialize_json_array(const std::vector & vec); + + std::string conn_info_; + /// conn_info_ in libpq key='value' form, with the default timeouts. Holds the password. + std::string connect_string_; + /// Password from conn_info_ or PGPASSWORD, removed from error text. Empty when unknown. + std::string password_; + /// No connection attempt before this time; set after a failed round. Guarded by mutex_. + mutable std::chrono::steady_clock::time_point next_connect_attempt_{}; + /// Reason of the last failed round, reported until the next one. Guarded by mutex_. + mutable std::string last_connect_error_; + /// The last connection round failed; the next one makes a single attempt. Guarded by mutex_. + mutable bool connect_failed_{false}; + /// Mutable: a reconnect replaces the connection but does not change the stored data. + mutable std::unique_ptr db_conn_; + /// Whether the schema was created on the current connection. Guarded by mutex_. + mutable bool schema_ready_{false}; + mutable std::mutex mutex_; + DebounceConfig config_; + size_t max_snapshots_per_fault_{0}; ///< 0 = unlimited + size_t max_near_misses_per_fault_{0}; ///< 0 = unlimited + bool retain_snapshots_on_clear_{false}; + /// Defaults to 1, the pre-#620 behaviour: a new recording replaces the old one. + /// 0 = unlimited, bounded only by max_total_storage_mb. + size_t max_rosbags_per_fault_{1}; + int max_retries_{1}; ///< Attempts after the first one. Must be >= 0. + unsigned reconnection_delay_{500}; +}; + +} // namespace ros2_medkit_fault_manager 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..07ad5203b 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 @@ -232,12 +232,19 @@ class SnapshotCapture { /// Mints one id per capture, shared by every row of that capture. /// - /// Seeded from storage at construction rather than started at zero: the ids - /// outlive the process on a persistent backend, and eviction protects the - /// HIGHEST id, so a counter that restarts below what is already stored makes the - /// freshly written capture the first one dropped. + /// Seeded from the highest id in storage: the ids outlive the process on a + /// persistent backend, and eviction protects the HIGHEST id, so a counter that + /// restarts below what is already stored makes the freshly written capture the + /// first one dropped. Seeded at construction, or on the first capture when the + /// store could not answer then; no capture runs unseeded. std::atomic capture_seq_{0}; + /// Seed capture_seq_ from storage once. False while the store is unreachable. + bool seed_capture_seq(); + + std::atomic capture_seq_seeded_{false}; + std::mutex capture_seq_seed_mutex_; + /// Compiled regex patterns (cached for performance) std::vector>> compiled_patterns_; 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..0da318ecc 100644 --- a/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp +++ b/src/ros2_medkit_fault_manager/src/fault_manager_node.cpp @@ -25,8 +25,12 @@ #include #include #include +#include #include "ros2_medkit_fault_manager/correlation/config_parser.hpp" +#ifdef POSTGRES_SUPPORT +#include "ros2_medkit_fault_manager/postgres_fault_storage.hpp" +#endif #include "ros2_medkit_fault_manager/sqlite_fault_storage.hpp" #include "ros2_medkit_fault_manager/time_utils.hpp" #include "ros2_medkit_msgs/msg/cluster_info.hpp" @@ -113,12 +117,18 @@ std::string validate_recording_id(const std::string & recording_id) { return ""; // Valid } +/// Response text for a request the fault storage server could not serve. +std::string storage_unavailable(const std::exception & e) { + return std::string("Fault storage unavailable: ") + e.what(); +} + } // namespace FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node("fault_manager", options) { // Declare and get parameters storage_type_ = declare_parameter("storage_type", "sqlite"); database_path_ = declare_parameter("database_path", "/var/lib/ros2_medkit/faults.db"); + database_url_ = declare_parameter("database_url", ""); auto confirmation_threshold_param = declare_parameter("confirmation_threshold", -1); if (confirmation_threshold_param > 0) { @@ -242,7 +252,13 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // Apply near-miss retention bound to storage (0 = unlimited). Applying it also trims a series // that a previous run left over the bound, which deletes history for good, so say when it does. - const size_t evicted_near_misses = storage_->set_max_near_misses_per_fault(static_cast(max_near_misses)); + // Storage that is unreachable at startup skips the trim; the bound still applies to new entries. + size_t evicted_near_misses = 0; + try { + evicted_near_misses = storage_->set_max_near_misses_per_fault(static_cast(max_near_misses)); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Skipped the startup near-miss trim, fault storage unavailable: %s", e.what()); + } if (evicted_near_misses > 0) { RCLCPP_WARN(get_logger(), "near_miss.max_per_fault=%ld dropped %zu stored near-miss entries that exceeded the bound. " @@ -270,20 +286,25 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // would behave inconsistently under the latch, so reclassify it as CLEARED at startup. Each flipped // fault is audited (the audit log is already constructed above); without this the reclassification // would be invisible to the audit log's verify(). + // Storage that is unreachable at startup skips it until the next start. if (!global_config_.healing_enabled) { - const auto reclassified = storage_->reclassify_healed_as_cleared(); - 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); - if (fault) { - audit_transition(kTransitionCleared, *fault, "startup_reclassify", reclassified_at_ns); + try { + const auto reclassified = storage_->reclassify_healed_as_cleared(); + 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); + if (fault) { + audit_transition(kTransitionCleared, *fault, "startup_reclassify", reclassified_at_ns); + } } } + RCLCPP_INFO(get_logger(), "Healing disabled: reclassified %zu stale HEALED fault(s) as CLEARED", + reclassified.size()); } - RCLCPP_INFO(get_logger(), "Healing disabled: reclassified %zu stale HEALED fault(s) as CLEARED", - reclassified.size()); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Startup HEALED reclassification not done, fault storage unavailable: %s", e.what()); } } @@ -447,7 +468,13 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // Create auto-confirmation timer if enabled if (auto_confirm_after_sec_ > 0.0) { auto_confirm_timer_ = create_wall_timer(std::chrono::seconds(1), [this]() { - const auto confirmed = storage_->check_time_based_confirmation(get_wall_clock_time()); + std::vector confirmed = {}; + try { + confirmed = storage_->check_time_based_confirmation(get_wall_clock_time()); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } if (confirmed.empty()) { return; } @@ -455,7 +482,13 @@ FaultManagerNode::FaultManagerNode(const rclcpp::NodeOptions & options) : Node(" // 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); + std::optional fault; + try { + fault = storage_->get_fault(fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } if (fault) { audit_transition(kTransitionConfirmed, *fault, "auto_confirm_timer", confirmed_at_ns); // A timer-driven confirmation is a confirmation: it has to reach the @@ -549,6 +582,21 @@ std::unique_ptr FaultManagerNode::create_storage() { return std::make_unique(database_path_); } +#ifdef POSTGRES_SUPPORT + if (storage_type_ == "postgres") { + auto postgres_fault_storage = std::make_unique(database_url_); + RCLCPP_INFO(get_logger(), "Using PostgreSQL fault storage (%s)", postgres_fault_storage->target().c_str()); + if (!postgres_fault_storage->connected()) { + RCLCPP_WARN(get_logger(), "PostgreSQL is unreachable: running without fault storage until it connects"); + } + return postgres_fault_storage; + } +#else + if (storage_type_ == "postgres") { + throw std::runtime_error("storage_type 'postgres' needs a build with -DPOSTGRES_SUPPORT=ON"); + } +#endif + RCLCPP_ERROR(get_logger(), "Unknown storage_type '%s', falling back to in-memory", storage_type_.c_str()); return std::make_unique(); } @@ -588,7 +636,7 @@ std::unique_ptr FaultManagerNode::create_audit_log() { } if (audit_path.empty()) { - if (database_path_ == ":memory:" || storage_type_ != "sqlite") { + if (database_path_ == ":memory:" || (storage_type_ != "sqlite" && storage_type_ != "postgres")) { audit_path = ":memory:"; } else { std::filesystem::path base(database_path_); @@ -804,7 +852,13 @@ void FaultManagerNode::handle_report_fault( } // Get status before update (if fault exists) - auto fault_before = storage_->get_fault(request->fault_code); + std::optional fault_before; + try { + fault_before = storage_->get_fault(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } std::string status_before = fault_before ? fault_before->status : ""; // Resolve per-entity debounce config (longest-prefix match on source_id) @@ -813,13 +867,25 @@ void FaultManagerNode::handle_report_fault( // Report the fault event (use wall clock time, not sim time, for proper timestamps) const rclcpp::Time event_time = get_wall_clock_time(); - bool is_new = storage_->report_fault_event(request->fault_code, request->event_type, request->severity, - request->description, request->source_id, event_time, resolved_config); + bool is_new = false; + try { + is_new = storage_->report_fault_event(request->fault_code, request->event_type, request->severity, + request->description, request->source_id, event_time, resolved_config); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } response->accepted = true; // Get updated fault state to publish event - auto fault_after = storage_->get_fault(request->fault_code); + std::optional fault_after; + try { + fault_after = storage_->get_fault(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } if (fault_after) { // Process through correlation engine (if enabled) // Only process FAILED events with correlation @@ -929,7 +995,12 @@ void FaultManagerNode::handle_report_fault( void FaultManagerNode::handle_list_faults( const std::shared_ptr & request, const std::shared_ptr & response) { - response->faults = storage_->list_faults(request->filter_by_severity, request->severity, request->statuses); + try { + response->faults = storage_->list_faults(request->filter_by_severity, request->severity, request->statuses); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + return; + } // Include correlation data if engine is enabled if (correlation_engine_) { @@ -1006,7 +1077,18 @@ void FaultManagerNode::handle_clear_fault( return; } - // Process through correlation engine first (to get auto-clear list). + // Storage first: a clear that the storage refuses must not change the correlation state. + bool cleared = false; + try { + cleared = storage_->clear_fault(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->success = false; + response->message = storage_unavailable(e); + return; + } + + // Then the correlation engine (to get the auto-clear list). // `skip_correlation_auto_clear` lets the caller opt out of cascade-clearing // correlated symptom fault codes. Per-entity DELETE routes set it to true // so they cannot reach across entity boundaries via the correlation graph. @@ -1016,8 +1098,6 @@ void FaultManagerNode::handle_clear_fault( auto_cleared_codes = clear_result.auto_cleared_codes; } - bool cleared = storage_->clear_fault(request->fault_code); - response->success = cleared; if (cleared) { // Evict cooldown tracking for cleared fault and auto-cleared symptoms @@ -1031,9 +1111,22 @@ void FaultManagerNode::handle_clear_fault( // Auto-clear correlated symptoms for (const auto & symptom_code : auto_cleared_codes) { - storage_->clear_fault(symptom_code); + try { + storage_->clear_fault(symptom_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->message = storage_unavailable(e); + return; + } if (audit_log_) { - auto symptom = storage_->get_fault(symptom_code); + std::optional symptom; + try { + symptom = storage_->get_fault(symptom_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->message = storage_unavailable(e); + return; + } if (symptom) { audit_transition(kTransitionCleared, *symptom, "clear_service", get_wall_clock_time().nanoseconds()); } @@ -1062,7 +1155,14 @@ void FaultManagerNode::handle_clear_fault( } // Publish EVENT_CLEARED - get the cleared fault to include in event - auto fault = storage_->get_fault(request->fault_code); + std::optional fault; + try { + fault = storage_->get_fault(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->message = storage_unavailable(e); + return; + } 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()); @@ -1098,7 +1198,15 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptrget_fault(request->fault_code); + std::optional fault; + try { + fault = storage_->get_fault(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->success = false; + response->error_message = storage_unavailable(e); + return; + } if (!fault) { response->success = false; response->error_message = "Fault not found: " + request->fault_code; @@ -1115,7 +1223,15 @@ 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); + std::vector stored_snapshots = {}; + try { + stored_snapshots = storage_->get_snapshots(request->fault_code); + } catch (const FaultStorage::IgnorableConnectionException & e) { + RCLCPP_WARN(get_logger(), "Failed to connect with the fault storage server: %s", e.what()); + response->success = false; + response->error_message = storage_unavailable(e); + return; + } 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 +1252,15 @@ void FaultManagerNode::handle_get_fault(const std::shared_ptr