diff --git a/waybionic_diagnostics/package.xml b/waybionic_diagnostics/package.xml new file mode 100644 index 0000000..591bc70 --- /dev/null +++ b/waybionic_diagnostics/package.xml @@ -0,0 +1,21 @@ + + + waybionic_diagnostics + 0.1.0 + CLI diagnostics tool for Waybionic ROS 2 systems. + + Aj Khidri + Apache-2.0 + + rclpy + diagnostic_msgs + + ament_copyright + ament_flake8 + ament_pep257 + python3-pytest + + + ament_python + + diff --git a/waybionic_diagnostics/resource/waybionic_diagnostics b/waybionic_diagnostics/resource/waybionic_diagnostics new file mode 100644 index 0000000..e69de29 diff --git a/waybionic_diagnostics/setup.cfg b/waybionic_diagnostics/setup.cfg new file mode 100644 index 0000000..f302dea --- /dev/null +++ b/waybionic_diagnostics/setup.cfg @@ -0,0 +1,5 @@ +[develop] +script_dir=$base/lib/waybionic_diagnostics + +[install] +install_scripts=$base/lib/waybionic_diagnostics diff --git a/waybionic_diagnostics/setup.py b/waybionic_diagnostics/setup.py new file mode 100644 index 0000000..8ee8b4b --- /dev/null +++ b/waybionic_diagnostics/setup.py @@ -0,0 +1,28 @@ +from setuptools import setup + +package_name = 'waybionic_diagnostics' + +setup( + name=package_name, + version='0.1.0', + packages=[package_name], + data_files=[ + ( + 'share/ament_index/resource_index/packages', + ['resource/' + package_name], + ), + ('share/' + package_name, ['package.xml']), + ], + install_requires=['setuptools'], + zip_safe=True, + maintainer='Aj Khidri', + maintainer_email='khidri.ajmal@gmail.com', + description='CLI diagnostics tool for Waybionic ROS 2 systems.', + license='Apache-2.0', + tests_require=['pytest'], + entry_points={ + 'console_scripts': [ + 'diagnostics = waybionic_diagnostics.cli:main', + ], + }, +) diff --git a/waybionic_diagnostics/test/test_cli.py b/waybionic_diagnostics/test/test_cli.py new file mode 100644 index 0000000..acc7c7a --- /dev/null +++ b/waybionic_diagnostics/test/test_cli.py @@ -0,0 +1,228 @@ +import time + +from diagnostic_msgs.msg import DiagnosticStatus, KeyValue + +from waybionic_diagnostics.cli import ( + DiagnosticsCliNode, + extract_value_and_unit, + overall_status, + print_snapshot, + status_name, +) + + +def test_status_name_mapping(): + assert status_name(DiagnosticStatus.OK) == 'OK' + assert status_name(DiagnosticStatus.WARN) == 'WARN' + assert status_name(DiagnosticStatus.ERROR) == 'FAULT' + assert status_name(DiagnosticStatus.STALE) == 'STALE' + assert status_name(99) == 'WARN' + + +def test_extract_value_and_unit(): + status = DiagnosticStatus() + status.values = [ + KeyValue(key='value', value='42.5'), + KeyValue(key='unit', value='C'), + ] + + value, unit = extract_value_and_unit(status) + + assert value == '42.5' + assert unit == 'C' + + +def test_extract_value_and_unit_case_insensitive(): + status = DiagnosticStatus() + status.values = [ + KeyValue(key='VALUE', value='0.85'), + KeyValue(key='UNIT', value='A'), + ] + + value, unit = extract_value_and_unit(status) + + assert value == '0.85' + assert unit == 'A' + + +def test_extract_value_and_unit_fallback(): + status = DiagnosticStatus() + status.values = [ + KeyValue(key='temperature', value='42.0'), + ] + + value, unit = extract_value_and_unit(status) + + assert value == '42.0' + assert unit == 'temperature' + + +def test_extract_value_and_unit_empty(): + status = DiagnosticStatus() + + value, unit = extract_value_and_unit(status) + + assert value is None + assert unit is None + + +def test_overall_status(): + assert overall_status([]) == 'STALE' + assert overall_status(['OK']) == 'OK' + assert overall_status(['OK', 'WARN']) == 'WARN' + assert overall_status(['OK', 'STALE']) == 'STALE' + assert overall_status(['STALE', 'FAULT']) == 'FAULT' + assert overall_status(['WARN', 'FAULT', 'STALE']) == 'FAULT' + + +def test_diagnostics_callback_updates_state(): + import rclpy + from diagnostic_msgs.msg import DiagnosticArray + + rclpy.init() + + node = DiagnosticsCliNode('/test_diagnostics') + + try: + message = DiagnosticArray() + status = DiagnosticStatus() + status.name = 'board.temperature' + status.level = DiagnosticStatus.OK + status.message = 'Normal Temperature' + status.values = [ + KeyValue(key='value', value='42.0'), + KeyValue(key='unit', value='C'), + ] + message.status = [status] + + before = time.monotonic() + node.diagnostics_callback(message) + + assert len(node.latest) == 1 + assert node.latest[0].name == 'board.temperature' + assert node.latest[0].level == DiagnosticStatus.OK + assert node.last_received_monotonic is not None + assert node.last_received_monotonic >= before + assert node.data_age() is not None + assert node.data_age() >= 0.0 + finally: + node.destroy_node() + rclpy.shutdown() + + +def test_diagnostics_callback_merges_entries_from_multiple_publishers(): + import rclpy + from diagnostic_msgs.msg import DiagnosticArray + + rclpy.init() + node = DiagnosticsCliNode('/test_diagnostics') + + try: + can_message = DiagnosticArray() + can_status = DiagnosticStatus() + can_status.name = 'can.bus' + can_status.level = DiagnosticStatus.ERROR + can_message.status = [can_status] + + imu_message = DiagnosticArray() + imu_status = DiagnosticStatus() + imu_status.name = 'imu.orientation' + imu_status.level = DiagnosticStatus.OK + imu_message.status = [imu_status] + + node.diagnostics_callback(can_message) + node.diagnostics_callback(imu_message) + + statuses = {status.name: status.level for status in node.latest} + assert statuses == { + 'can.bus': DiagnosticStatus.ERROR, + 'imu.orientation': DiagnosticStatus.OK, + } + finally: + node.destroy_node() + rclpy.shutdown() + + +def test_run_snapshot_collects_queued_publishers(monkeypatch, capsys): + import rclpy + from diagnostic_msgs.msg import DiagnosticArray + from waybionic_diagnostics import cli + + rclpy.init() + node = DiagnosticsCliNode('/test_diagnostics') + can_message = DiagnosticArray() + can_status = DiagnosticStatus() + can_status.name = 'can.bus' + can_status.level = DiagnosticStatus.ERROR + can_message.status = [can_status] + imu_message = DiagnosticArray() + imu_status = DiagnosticStatus() + imu_status.name = 'imu.orientation' + imu_status.level = DiagnosticStatus.OK + imu_message.status = [imu_status] + + try: + node.diagnostics_callback(imu_message) + node.diagnostics_callback(can_message) + monkeypatch.setattr(cli, 'SNAPSHOT_COLLECTION_SECONDS', 0.0) + cli.run_snapshot(node) + + output = capsys.readouterr().out + assert 'can.bus' in output + assert 'Overall: FAULT' in output + finally: + node.destroy_node() + rclpy.shutdown() + + +def test_diagnostics_callback_uses_header_stamp_for_sample_age(monkeypatch, capsys): + import rclpy + from diagnostic_msgs.msg import DiagnosticArray + + rclpy.init() + node = DiagnosticsCliNode('/test_diagnostics') + + try: + message = DiagnosticArray() + message.header.stamp.sec = 60 + status = DiagnosticStatus() + status.name = 'board.temperature' + status.level = DiagnosticStatus.OK + message.status = [status] + + monkeypatch.setattr('waybionic_diagnostics.cli.time.time', lambda: 70.0) + node.diagnostics_callback(message) + + assert node.status_age(status) == 10.0 + print_snapshot(node) + assert 'STALE' in capsys.readouterr().out + finally: + node.destroy_node() + rclpy.shutdown() + + +def test_main_forwards_ros_args_to_rclpy_init(monkeypatch): + from waybionic_diagnostics import cli + + captured = {} + + class FakeNode: + def __init__(self, topic): + captured['topic'] = topic + + def destroy_node(self): + captured['destroyed'] = True + + monkeypatch.setattr(cli, 'DiagnosticsCliNode', FakeNode) + monkeypatch.setattr(cli, 'run_snapshot', lambda node: captured.setdefault('snapshot', node)) + monkeypatch.setattr(cli.rclpy, 'init', lambda args=None: captured.setdefault('init_args', args)) + monkeypatch.setattr(cli.rclpy, 'ok', lambda: False) + monkeypatch.setattr(cli.rclpy, 'shutdown', lambda: captured.setdefault('shutdown', True)) + + cli.main(['--topic', '/diagnostics_custom', '--ros-args', '-r', '__node:=renamed']) + + assert captured['topic'] == '/diagnostics_custom' + assert captured['init_args'] == ['--ros-args', '-r', '__node:=renamed'] + assert 'snapshot' in captured + assert captured['destroyed'] is True + assert 'shutdown' not in captured diff --git a/waybionic_diagnostics/waybionic_diagnostics/__init__.py b/waybionic_diagnostics/waybionic_diagnostics/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/waybionic_diagnostics/waybionic_diagnostics/cli.py b/waybionic_diagnostics/waybionic_diagnostics/cli.py new file mode 100644 index 0000000..227a1ea --- /dev/null +++ b/waybionic_diagnostics/waybionic_diagnostics/cli.py @@ -0,0 +1,278 @@ +import argparse +import sys +import time +from typing import Dict, List, Optional + +import rclpy +from diagnostic_msgs.msg import DiagnosticArray, DiagnosticStatus +from rclpy.node import Node + + +STALE_AFTER_SECONDS = 5.0 +SNAPSHOT_COLLECTION_SECONDS = 0.1 + + +def status_name(level: int) -> str: + """Convert a ROS diagnostic level into the CLI status name.""" + if level == DiagnosticStatus.OK: + return 'OK' + if level == DiagnosticStatus.WARN: + return 'WARN' + if level == DiagnosticStatus.ERROR: + return 'FAULT' + if level == DiagnosticStatus.STALE: + return 'STALE' + return 'WARN' + + +def extract_value_and_unit(status: DiagnosticStatus): + """Extract value/unit using the same rules as the RViz diagnostics source.""" + value: Optional[str] = None + unit: Optional[str] = None + + for key_value in status.values: + key = key_value.key.lower() + + if key == 'value': + value = key_value.value + continue + + if key == 'unit': + unit = key_value.value + continue + + if value is None and key_value.value: + value = key_value.value + unit = key_value.key + + return value, unit + + +class DiagnosticsCliNode(Node): + """ROS 2 subscriber and state holder for the diagnostics CLI.""" + + def __init__(self, topic: str): + super().__init__('waybionic_diagnostics_cli') + + self.topic = topic + self.latest: List[DiagnosticStatus] = [] + self._received_monotonic_by_name: Dict[str, float] = {} + self._sample_time_by_name: Dict[str, float] = {} + self.last_received_monotonic: Optional[float] = None + + self.subscription = self.create_subscription( + DiagnosticArray, + topic, + self.diagnostics_callback, + 10, + ) + + def diagnostics_callback(self, message: DiagnosticArray): + """Merge the publisher's entries into the diagnostic cache.""" + received_monotonic = time.monotonic() + received_time = time.time() + stamp = message.header.stamp + sample_time = stamp.sec + stamp.nanosec * 1e-9 + if sample_time == 0.0: + sample_time = received_time + + latest_by_name = {status.name: status for status in self.latest} + for status in message.status: + latest_by_name[status.name] = status + self._received_monotonic_by_name[status.name] = received_monotonic + self._sample_time_by_name[status.name] = sample_time + + self.latest = list(latest_by_name.values()) + self.last_received_monotonic = received_monotonic + + def data_age(self) -> Optional[float]: + """Return seconds since the most recent diagnostic message was received.""" + if self.last_received_monotonic is None: + return None + + return max(0.0, time.monotonic() - self.last_received_monotonic) + + def status_age(self, status: DiagnosticStatus) -> Optional[float]: + """Return the sample age for one cached diagnostic status.""" + sample_time = self._sample_time_by_name.get(status.name) + received_monotonic = self._received_monotonic_by_name.get(status.name) + if sample_time is None or received_monotonic is None: + return None + + sample_age = max(0.0, time.time() - sample_time) + receive_age = max(0.0, time.monotonic() - received_monotonic) + return max(sample_age, receive_age) + + +def overall_status(statuses: List[str]) -> str: + """Calculate the worst status in a diagnostic snapshot.""" + priority = { + 'OK': 0, + 'WARN': 1, + 'STALE': 2, + 'FAULT': 3, + } + + if not statuses: + return 'STALE' + + return max(statuses, key=lambda status: priority.get(status, 1)) + + +def clear_screen(): + """Clear the terminal for watch mode.""" + sys.stdout.write('\033[2J\033[H') + sys.stdout.flush() + + +def print_snapshot(node: DiagnosticsCliNode, clear: bool = False): + """Render the latest diagnostic state.""" + if clear: + clear_screen() + + print('WAYBIONIC DIAGNOSTICS') + print('─' * 100) + + if not node.latest: + age = node.data_age() + + if age is None: + print(f'Waiting for messages on {node.topic}...') + else: + print( + f'No diagnostic entries received. ' + f'Last message was {age:.1f}s ago.' + ) + + print('─' * 100) + print('Overall: STALE') + return + + print( + f'{"Signal":<32} ' + f'{"Status":<8} ' + f'{"Value":<16} ' + f'{"Age":<10} ' + f'Message' + ) + print('─' * 100) + + rendered_statuses = [] + + for status in node.latest: + normalized = status_name(status.level) + age = node.status_age(status) + sample_stale = age is not None and age > STALE_AFTER_SECONDS + + if sample_stale and normalized in ('OK', 'WARN'): + normalized = 'STALE' + + value, unit = extract_value_and_unit(status) + + if value is None: + display_value = '-' + elif unit: + display_value = f'{value} {unit}' + else: + display_value = value + + if age is None: + display_age = '-' + else: + display_age = f'{age:.1f}s' + + message = status.message or '-' + + if sample_stale and status.level in ( + DiagnosticStatus.OK, + DiagnosticStatus.WARN, + ): + message = f'No recent update from {node.topic}' + + print( + f'{status.name:<32.32} ' + f'{normalized:<8} ' + f'{display_value:<16.16} ' + f'{display_age:<10} ' + f'{message}' + ) + + rendered_statuses.append(normalized) + + print('─' * 100) + print(f'Overall: {overall_status(rendered_statuses)}') + + stream_age = node.data_age() + if stream_age is not None: + print(f'Diagnostics stream age: {stream_age:.1f}s') + + +def wait_for_first_message(node: DiagnosticsCliNode, timeout: float = 5.0): + """Wait briefly for the first diagnostics message.""" + deadline = time.monotonic() + timeout + + while rclpy.ok() and node.last_received_monotonic is None: + if time.monotonic() >= deadline: + return False + + rclpy.spin_once(node, timeout_sec=0.1) + + return node.last_received_monotonic is not None + + +def run_snapshot(node: DiagnosticsCliNode): + """Run one diagnostic snapshot.""" + if wait_for_first_message(node): + deadline = time.monotonic() + SNAPSHOT_COLLECTION_SECONDS + while rclpy.ok() and time.monotonic() < deadline: + rclpy.spin_once(node, timeout_sec=0.0) + print_snapshot(node) + + +def run_watch(node: DiagnosticsCliNode): + """Continuously display diagnostics.""" + try: + while rclpy.ok(): + rclpy.spin_once(node, timeout_sec=0.2) + print_snapshot(node, clear=True) + time.sleep(0.3) + except KeyboardInterrupt: + pass + + +def main(args=None): + """Run the Waybionic diagnostics CLI.""" + parser = argparse.ArgumentParser( + description='Waybionic ROS 2 diagnostics viewer.' + ) + + parser.add_argument( + '--topic', + default='/diagnostics', + help='ROS 2 DiagnosticArray topic (default: /diagnostics)', + ) + + parser.add_argument( + '--watch', + action='store_true', + help='Continuously monitor diagnostics.', + ) + + parsed_args, ros_args = parser.parse_known_args(args) + + rclpy.init(args=ros_args) + node = DiagnosticsCliNode(parsed_args.topic) + + try: + if parsed_args.watch: + run_watch(node) + else: + run_snapshot(node) + finally: + node.destroy_node() + if rclpy.ok(): + rclpy.shutdown() + + +if __name__ == '__main__': + main() diff --git a/waybionic_rviz_plugins/scripts/temporary_diagnostics_publisher.py b/waybionic_rviz_plugins/scripts/temporary_diagnostics_publisher.py index 4a3f7a3..59081e4 100755 --- a/waybionic_rviz_plugins/scripts/temporary_diagnostics_publisher.py +++ b/waybionic_rviz_plugins/scripts/temporary_diagnostics_publisher.py @@ -70,30 +70,35 @@ def build_statuses(self, mode): make_status( 'board.temperature', DiagnosticStatus.OK, + 'Normal Temperature', value=f'{42.0 + pulse:.1f}', unit='C', ), make_status( 'motor.current', DiagnosticStatus.OK, + 'Normal Motor Current', value=f'{0.8 + pulse * 0.05:.2f}', unit='A', ), make_status( 'imu.roll', DiagnosticStatus.OK, + 'Roll Normal', value=f'{1.2 + pulse * 0.1:.1f}', unit='deg', ), make_status( 'imu.pitch', DiagnosticStatus.OK, + 'Pitch Normal', value=f'{-0.4 + pulse * 0.1:.1f}', unit='deg', ), make_status( 'imu.yaw', DiagnosticStatus.OK, + 'Yaw Normal', value=f'{12.9 + pulse * 0.2:.1f}', unit='deg', ), diff --git a/waybionic_rviz_plugins/src/ros_diagnostics_source.cpp b/waybionic_rviz_plugins/src/ros_diagnostics_source.cpp index 33a3387..8948cda 100644 --- a/waybionic_rviz_plugins/src/ros_diagnostics_source.cpp +++ b/waybionic_rviz_plugins/src/ros_diagnostics_source.cpp @@ -71,7 +71,7 @@ DiagnosticMessage toDiagnosticMessage( const auto normalized_status = mapLevel(status.level); std::optional alert_message; - if (hasContent(status.message) && normalized_status != DiagnosticStatus::Ok) { + if (normalized_status != DiagnosticStatus::Ok && hasContent(status.message)) { alert_message = status.message; } diff --git a/waybionic_rviz_plugins/test/test_ros_diagnostics_source.cpp b/waybionic_rviz_plugins/test/test_ros_diagnostics_source.cpp index 63448e0..7de6b8d 100644 --- a/waybionic_rviz_plugins/test/test_ros_diagnostics_source.cpp +++ b/waybionic_rviz_plugins/test/test_ros_diagnostics_source.cpp @@ -175,6 +175,22 @@ TEST_F(DiagnosticsTrafficFixture, NormalizesReceivedDiagnostics) stopTraffic(); } +TEST_F(DiagnosticsTrafficFixture, OkDiagnosticsDoNotCarryAlertMessage) +{ + const auto source = makeSource(); + + publisher_->publish(makeArray(diagnostic_msgs::msg::DiagnosticStatus::OK, "42")); + ASSERT_TRUE(waitFor([&]() { + const auto messages = source->messages(now()); + return messages.size() == 1u && messages.front().signal_name == "board.temperature"; + }, 5s)); + + const auto messages = source->messages(now()); + ASSERT_EQ(messages.size(), 1u); + EXPECT_EQ(messages.front().status, DiagnosticStatus::Ok); + EXPECT_FALSE(messages.front().alert_message.has_value()); +} + TEST_F(DiagnosticsTrafficFixture, StopFreezesStateAndIgnoresLaterMessages) { auto source = makeSource();