From 417c64f697dc6f72890471f4ce52dec048af4276 Mon Sep 17 00:00:00 2001 From: Gabriela Vitez Date: Tue, 11 Aug 2026 11:19:04 +0200 Subject: [PATCH 1/4] feat: convert SetGoal service to action --- foreman/foreman/adapters/__init__.py | 4 +- .../adapters/ros_set_goal_action_server.py | 108 ++++++++++ .../foreman/adapters/ros_set_goal_server.py | 44 ---- foreman/foreman/node.py | 2 +- .../test/test_ros_set_goal_action_server.py | 189 ++++++++++++++++++ foreman_msgs/CMakeLists.txt | 13 +- foreman_msgs/action/SetGoal.action | 11 + foreman_msgs/msg/ForemanErrorState.msg | 9 + foreman_msgs/package.xml | 4 +- foreman_msgs/srv/SetGoal.srv | 6 - 10 files changed, 333 insertions(+), 57 deletions(-) create mode 100644 foreman/foreman/adapters/ros_set_goal_action_server.py delete mode 100644 foreman/foreman/adapters/ros_set_goal_server.py create mode 100644 foreman/test/test_ros_set_goal_action_server.py create mode 100644 foreman_msgs/action/SetGoal.action create mode 100644 foreman_msgs/msg/ForemanErrorState.msg delete mode 100644 foreman_msgs/srv/SetGoal.srv diff --git a/foreman/foreman/adapters/__init__.py b/foreman/foreman/adapters/__init__.py index 9485224..57af74e 100644 --- a/foreman/foreman/adapters/__init__.py +++ b/foreman/foreman/adapters/__init__.py @@ -2,7 +2,7 @@ from .controller_manager_service_caller import ControllerManagerServiceCaller from .lifecycle_node_service_caller import LifecycleNodeServiceCaller from .ros_node_parameters import RosNodeParameters -from .ros_set_goal_server import RosSetGoalServer +from .ros_set_goal_action_server import RosSetGoalActionServer try: from .datalayer.datalayer_adapter import DatalayerAdapter @@ -15,7 +15,7 @@ "ComponentStateMonitor", "ControllerManagerServiceCaller", "LifecycleNodeServiceCaller", - "RosSetGoalServer", + "RosSetGoalActionServer", "RosNodeParameters", "DatalayerAdapter" ] diff --git a/foreman/foreman/adapters/ros_set_goal_action_server.py b/foreman/foreman/adapters/ros_set_goal_action_server.py new file mode 100644 index 0000000..1a5a367 --- /dev/null +++ b/foreman/foreman/adapters/ros_set_goal_action_server.py @@ -0,0 +1,108 @@ +import time + +from rclpy.action import ActionServer +from rclpy.action import CancelResponse +from rclpy.action import GoalResponse +from rclpy.node import Node + +from foreman.engine import ForemanEngine +from foreman.types import ErrorSnapshot +from foreman_msgs.action import SetGoal +from foreman_msgs.msg import ForemanErrorState + + +def _to_error_msg(snapshot: ErrorSnapshot) -> ForemanErrorState: + """Convert the engine's error snapshot into its ROS representation.""" + msg = ForemanErrorState() + msg.is_error = snapshot.is_error + msg.category = snapshot.category + msg.message = snapshot.message + msg.components = list(snapshot.components or []) + return msg + + +class RosSetGoalActionServer: + """Drive the system to a named goal with a ROS 2 action.""" + + def __init__(self, node: Node, engine: ForemanEngine, poll_period: float = 0.05): + self._node = node + self._engine = engine + self._poll_period = poll_period + self.logger_prefix = "Adapters.RosSetGoalActionServer:" + + self._action_server = ActionServer( + node, + SetGoal, + 'foreman/set_goal', + execute_callback=self._execute, + goal_callback=self._on_goal_request, + cancel_callback=self._on_cancel_request, + callback_group=node.callback_group_subscriber + ) + + self._node.get_logger().info(f"{self.logger_prefix} Action /foreman/set_goal is ready.") + + def _on_goal_request(self, goal_request) -> GoalResponse: + self._node.get_logger().info( + f"{self.logger_prefix} Received request for goal '{goal_request.goal}'") + return GoalResponse.ACCEPT + + def _on_cancel_request(self, goal_handle) -> CancelResponse: + """Cancel waiting without rolling the system back.""" + del goal_handle + return CancelResponse.ACCEPT + + def _execute(self, goal_handle): + goal_name = goal_handle.request.goal + result = SetGoal.Result() + + engine_response = self._engine.request_goal(goal_name) + if not engine_response.success: + self._node.get_logger().warning(f"{engine_response.message}") + result.success = False + result.message = engine_response.message + result.error = _to_error_msg(self._engine.get_engine_snapshot().error) + goal_handle.abort() + return result + + self._node.get_logger().info(f"{engine_response.message}") + + feedback = SetGoal.Feedback() + while True: + if not goal_handle.is_active: + result.success = False + result.message = f"Goal '{goal_name}' was preempted." + return result + + if goal_handle.is_cancel_requested: + result.success = False + result.message = f"Stopped waiting for goal '{goal_name}'." + result.error = _to_error_msg(self._engine.get_engine_snapshot().error) + goal_handle.canceled() + return result + + snapshot = self._engine.get_engine_snapshot() + error_msg = _to_error_msg(snapshot.error) + + if snapshot.error.is_error: + result.success = False + result.message = f"[{snapshot.error.category}] {snapshot.error.message}" + result.error = error_msg + self._node.get_logger().error( + f"{self.logger_prefix} Goal '{goal_name}' aborted: {result.message}") + goal_handle.abort() + return result + + if snapshot.at_goal: + result.success = True + result.message = f"Goal '{goal_name}' reached." + result.error = error_msg + self._node.get_logger().info(f"{self.logger_prefix} {result.message}") + goal_handle.succeed() + return result + + feedback.at_goal = False + feedback.error = error_msg + goal_handle.publish_feedback(feedback) + + time.sleep(self._poll_period) diff --git a/foreman/foreman/adapters/ros_set_goal_server.py b/foreman/foreman/adapters/ros_set_goal_server.py deleted file mode 100644 index 1384c2c..0000000 --- a/foreman/foreman/adapters/ros_set_goal_server.py +++ /dev/null @@ -1,44 +0,0 @@ -from rclpy.node import Node - -from foreman.engine import ForemanEngine -from foreman_msgs.srv import SetGoal - - -class RosSetGoalServer: - """ROS 2 service to set a named goal for Foreman Engine.""" - - def __init__(self, node: Node, engine: ForemanEngine): - self._node = node - self._engine = engine - self.logger_prefix = "Adapters.RosSetGoalServer:" - # Using MutuallyExclusiveCallbackGroup - # If a service is processing, we reject new service requests. - self._srv = self._node.create_service( - SetGoal, - 'foreman/set_goal', - self._handle_set_goal, - callback_group=self._node.callback_group_services - ) - - print() - - self._node.get_logger().info(f"{self.logger_prefix} Service /foreman/set_goal is ready.") - - def _handle_set_goal(self, request, response): - """Set the target system state.""" - goal_name = request.goal - # TODO: demote some of these to DEBUG logs. - self._node.get_logger().info( - f"{self.logger_prefix} Received request for goal '{goal_name}'") - - engine_response = self._engine.request_goal(goal_name) - - response.success = engine_response.success - response.message = engine_response.message - - if not engine_response.success: - self._node.get_logger().warning(f"{engine_response.message}") - else: - self._node.get_logger().info(f"{engine_response.message}") - - return response diff --git a/foreman/foreman/node.py b/foreman/foreman/node.py index 0194ef0..cb330ee 100644 --- a/foreman/foreman/node.py +++ b/foreman/foreman/node.py @@ -56,7 +56,7 @@ def __init__(self): lifecycle_nodes=self.foreman_config.lifecycle_nodes ) - self.ros_set_goal_server = adapters.RosSetGoalServer( + self.ros_set_goal_server = adapters.RosSetGoalActionServer( node=self, engine=self.foreman_engine ) diff --git a/foreman/test/test_ros_set_goal_action_server.py b/foreman/test/test_ros_set_goal_action_server.py new file mode 100644 index 0000000..7d10592 --- /dev/null +++ b/foreman/test/test_ros_set_goal_action_server.py @@ -0,0 +1,189 @@ +import unittest +from unittest.mock import MagicMock + +import rclpy +from rclpy.action import get_action_names_and_types +from rclpy.callback_groups import ReentrantCallbackGroup + +from foreman.adapters.ros_set_goal_action_server import _to_error_msg +from foreman.adapters.ros_set_goal_action_server import RosSetGoalActionServer +from foreman.types import ErrorSnapshot +from foreman.types import ForemanErrorCategory +from foreman.types import ForemanResponse +from foreman.types import ForemanSnapshot + + +def _snapshot(goal="force_ctrl", ready=True, at_goal=False, error=None): + """Build a ForemanSnapshot with a no-error default.""" + if error is None: + error = ErrorSnapshot( + is_error=False, + category=ForemanErrorCategory.NONE.value, + message="", + components=[] + ) + return ForemanSnapshot(goal=goal, ready=ready, at_goal=at_goal, error=error, components=[]) + + +def _error_snapshot(category=ForemanErrorCategory.EXECUTION, message="boom", components=None): + return ErrorSnapshot( + is_error=True, + category=category.value, + message=message, + components=components if components is not None else ["ctrl_a"] + ) + + +def _goal_handle(goal_name="force_ctrl"): + """Fake action goal handle: active, not cancelled, records terminal calls.""" + handle = MagicMock() + handle.request.goal = goal_name + handle.is_active = True + handle.is_cancel_requested = False + return handle + + +class TestToErrorMsg(unittest.TestCase): + def test_maps_every_field(self): + msg = _to_error_msg(_error_snapshot(components=["a", "b"])) + self.assertTrue(msg.is_error) + self.assertEqual(msg.category, ForemanErrorCategory.EXECUTION.value) + self.assertEqual(msg.message, "boom") + self.assertEqual(list(msg.components), ["a", "b"]) + + def test_none_components_become_empty_list(self): + msg = _to_error_msg(ErrorSnapshot( + is_error=False, category="None", message="", components=None)) + self.assertEqual(list(msg.components), []) + + +class TestRosSetGoalActionServer(unittest.TestCase): + @classmethod + def setUpClass(cls): + rclpy.init() + + @classmethod + def tearDownClass(cls): + rclpy.shutdown() + + def setUp(self): + self.node = rclpy.create_node("test_set_goal_action") + # the adapter reads this; the real node sets it, a bare node doesn't + self.node.callback_group_subscriber = ReentrantCallbackGroup() + self.addCleanup(self.node.destroy_node) + self.engine = MagicMock() + + def _server(self): + # poll_period=0 so the wait loop does not slow the tests down + return RosSetGoalActionServer(self.node, self.engine, poll_period=0.0) + + def test_action_is_advertised_as_foreman_set_goal(self): + self._server() + advertised = dict(get_action_names_and_types(node=self.node)) + self.assertIn("/foreman/set_goal", advertised) + self.assertEqual(advertised["/foreman/set_goal"], ["foreman_msgs/action/SetGoal"]) + + def test_rejected_goal_aborts_without_waiting(self): + self.engine.request_goal.return_value = ForemanResponse( + False, "Goal 'nope' not found in configuration.") + self.engine.get_engine_snapshot.return_value = _snapshot() + handle = _goal_handle("nope") + + result = self._server()._execute(handle) + + handle.abort.assert_called_once() + handle.succeed.assert_not_called() + self.assertFalse(result.success) + self.assertIn("not found", result.message) + + def test_succeeds_only_once_engine_reports_at_goal(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + # two polls in transition, then arrived + self.engine.get_engine_snapshot.side_effect = [ + _snapshot(at_goal=False), + _snapshot(at_goal=False), + _snapshot(at_goal=True), + ] + handle = _goal_handle() + + result = self._server()._execute(handle) + + handle.succeed.assert_called_once() + handle.abort.assert_not_called() + self.assertTrue(result.success) + # feedback published for each poll that was still in transition + self.assertEqual(handle.publish_feedback.call_count, 2) + + def test_feedback_carries_current_error_state(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.side_effect = [ + _snapshot(at_goal=False), + _snapshot(at_goal=True), + ] + handle = _goal_handle() + + self._server()._execute(handle) + + feedback = handle.publish_feedback.call_args[0][0] + self.assertFalse(feedback.at_goal) + self.assertFalse(feedback.error.is_error) + + def test_engine_error_aborts_and_reports_it(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.side_effect = [ + _snapshot(at_goal=False), + _snapshot(error=_error_snapshot(message="Service rejected the transition.")), + ] + handle = _goal_handle() + + result = self._server()._execute(handle) + + handle.abort.assert_called_once() + handle.succeed.assert_not_called() + self.assertFalse(result.success) + self.assertIn("Service rejected the transition.", result.message) + self.assertTrue(result.error.is_error) + self.assertEqual(result.error.category, ForemanErrorCategory.EXECUTION.value) + self.assertEqual(list(result.error.components), ["ctrl_a"]) + + def test_cancel_stops_waiting(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=False) + handle = _goal_handle() + handle.is_cancel_requested = True + + result = self._server()._execute(handle) + + handle.canceled.assert_called_once() + handle.succeed.assert_not_called() + handle.abort.assert_not_called() + self.assertFalse(result.success) + + def test_preempted_goal_returns_without_terminal_call(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=False) + handle = _goal_handle() + handle.is_active = False + + result = self._server()._execute(handle) + + handle.succeed.assert_not_called() + handle.abort.assert_not_called() + handle.canceled.assert_not_called() + self.assertFalse(result.success) + + def test_already_at_goal_succeeds_without_feedback(self): + self.engine.request_goal.return_value = ForemanResponse( + True, "Already at goal 'force_ctrl'.") + self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=True) + handle = _goal_handle() + + result = self._server()._execute(handle) + + handle.succeed.assert_called_once() + handle.publish_feedback.assert_not_called() + self.assertTrue(result.success) + + +if __name__ == "__main__": + unittest.main() diff --git a/foreman_msgs/CMakeLists.txt b/foreman_msgs/CMakeLists.txt index eaf35a1..57495a3 100644 --- a/foreman_msgs/CMakeLists.txt +++ b/foreman_msgs/CMakeLists.txt @@ -12,13 +12,20 @@ endif() find_package(ament_cmake REQUIRED) find_package(rosidl_default_generators REQUIRED) +find_package(action_msgs REQUIRED) -set(srv_files - "srv/SetGoal.srv" +set(msg_files + "msg/ForemanErrorState.msg" +) + +set(action_files + "action/SetGoal.action" ) rosidl_generate_interfaces(${PROJECT_NAME} - ${srv_files} + ${msg_files} + ${action_files} + DEPENDENCIES action_msgs ) ament_export_dependencies(rosidl_default_runtime) diff --git a/foreman_msgs/action/SetGoal.action b/foreman_msgs/action/SetGoal.action new file mode 100644 index 0000000..874b5a8 --- /dev/null +++ b/foreman_msgs/action/SetGoal.action @@ -0,0 +1,11 @@ +# Drive the system to a named goal state. +# Succeeds once the goal is reached, not when it is accepted. + +string goal +--- +bool success +string message +ForemanErrorState error +--- +bool at_goal +ForemanErrorState error diff --git a/foreman_msgs/msg/ForemanErrorState.msg b/foreman_msgs/msg/ForemanErrorState.msg new file mode 100644 index 0000000..ccceab7 --- /dev/null +++ b/foreman_msgs/msg/ForemanErrorState.msg @@ -0,0 +1,9 @@ +# Domain error reported by the Foreman engine. +# Mirrors foreman.types.ErrorSnapshot. + +bool is_error +# ForemanErrorCategory value: TransportError, ExecutionError, +# UnexpectedStateError, PlannerError, or None. +string category +string message +string[] components diff --git a/foreman_msgs/package.xml b/foreman_msgs/package.xml index 20ec77f..f91d231 100644 --- a/foreman_msgs/package.xml +++ b/foreman_msgs/package.xml @@ -3,7 +3,7 @@ foreman_msgs 0.0.1 - Message and service definitions for the Foreman system. + Message and action definitions for the Foreman system. Nikola Banovic Apache License 2.0 @@ -11,6 +11,8 @@ ament_cmake rosidl_default_generators + action_msgs + rosidl_default_runtime ament_lint_auto diff --git a/foreman_msgs/srv/SetGoal.srv b/foreman_msgs/srv/SetGoal.srv deleted file mode 100644 index fb05de4..0000000 --- a/foreman_msgs/srv/SetGoal.srv +++ /dev/null @@ -1,6 +0,0 @@ -# TODO: add description - -string goal ---- -bool success -string message From d500690be4adecc0bd9e72925de8226ad9a9f14e Mon Sep 17 00:00:00 2001 From: Gabriela Vitez Date: Wed, 12 Aug 2026 10:56:32 +0200 Subject: [PATCH 2/4] Return the service but make it blocking --- foreman/foreman/adapters/__init__.py | 2 + .../foreman/adapters/ros_set_goal_server.py | 78 ++++++++++++ foreman/foreman/node.py | 6 +- foreman/test/test_ros_set_goal_server.py | 115 ++++++++++++++++++ foreman_msgs/CMakeLists.txt | 5 + foreman_msgs/srv/SetGoal.srv | 6 + 6 files changed, 211 insertions(+), 1 deletion(-) create mode 100644 foreman/foreman/adapters/ros_set_goal_server.py create mode 100644 foreman/test/test_ros_set_goal_server.py create mode 100644 foreman_msgs/srv/SetGoal.srv diff --git a/foreman/foreman/adapters/__init__.py b/foreman/foreman/adapters/__init__.py index 57af74e..696a4f3 100644 --- a/foreman/foreman/adapters/__init__.py +++ b/foreman/foreman/adapters/__init__.py @@ -3,6 +3,7 @@ from .lifecycle_node_service_caller import LifecycleNodeServiceCaller from .ros_node_parameters import RosNodeParameters from .ros_set_goal_action_server import RosSetGoalActionServer +from .ros_set_goal_server import RosSetGoalServer try: from .datalayer.datalayer_adapter import DatalayerAdapter @@ -16,6 +17,7 @@ "ControllerManagerServiceCaller", "LifecycleNodeServiceCaller", "RosSetGoalActionServer", + "RosSetGoalServer", "RosNodeParameters", "DatalayerAdapter" ] diff --git a/foreman/foreman/adapters/ros_set_goal_server.py b/foreman/foreman/adapters/ros_set_goal_server.py new file mode 100644 index 0000000..bc6f6f7 --- /dev/null +++ b/foreman/foreman/adapters/ros_set_goal_server.py @@ -0,0 +1,78 @@ +import time + +from rclpy.callback_groups import MutuallyExclusiveCallbackGroup +from rclpy.node import Node + +from foreman.engine import ForemanEngine +from foreman_msgs.srv import SetGoal + + +class RosSetGoalServer: + """ROS 2 service to set a named goal for Foreman Engine.""" + + def __init__(self, node: Node, engine: ForemanEngine): + self._node = node + self._engine = engine + self._poll_period = 0.05 + self.logger_prefix = "Adapters.RosSetGoalServer:" + # Using MutuallyExclusiveCallbackGroup + # If a service is processing, we reject new service requests. + self._callback_group = MutuallyExclusiveCallbackGroup() + + self._srv = self._node.create_service( + SetGoal, + 'foreman/set_goal', + self._handle_set_goal, + callback_group=self._callback_group + ) + + print() + + self._node.get_logger().info(f"{self.logger_prefix} Service /foreman/set_goal is ready.") + + def _handle_set_goal(self, request, response): + """Set the target system state.""" + goal_name = request.goal + # TODO: demote some of these to DEBUG logs. + self._node.get_logger().info( + f"{self.logger_prefix} Received request for goal '{goal_name}'") + + engine_response = self._engine.request_goal(goal_name) + if not engine_response.success: + self._node.get_logger().warning(f"{engine_response.message}") + response.success = False + response.message = engine_response.message + return response + + self._node.get_logger().info(f"{engine_response.message}") + + while True: + snapshot = self._engine.get_engine_snapshot() + + if snapshot.error.is_error: + response.success = False + response.message = ( + f"[{snapshot.error.category}] {snapshot.error.message}" + ) + self._node.get_logger().error( + f"{self.logger_prefix} Goal '{goal_name}' aborted: " + f"{response.message}") + return response + + if snapshot.goal != goal_name: + response.success = False + response.message = ( + f"Goal '{goal_name}' was preempted by goal '{snapshot.goal}'." + ) + self._node.get_logger().warning( + f"{self.logger_prefix} {response.message}") + return response + + if snapshot.at_goal: + response.success = True + response.message = f"Goal '{goal_name}' reached." + self._node.get_logger().info( + f"{self.logger_prefix} {response.message}") + return response + + time.sleep(self._poll_period) diff --git a/foreman/foreman/node.py b/foreman/foreman/node.py index cb330ee..f0a89ac 100644 --- a/foreman/foreman/node.py +++ b/foreman/foreman/node.py @@ -56,7 +56,11 @@ def __init__(self): lifecycle_nodes=self.foreman_config.lifecycle_nodes ) - self.ros_set_goal_server = adapters.RosSetGoalActionServer( + self.ros_set_goal_action_server = adapters.RosSetGoalActionServer( + node=self, + engine=self.foreman_engine + ) + self.ros_set_goal_server = adapters.RosSetGoalServer( node=self, engine=self.foreman_engine ) diff --git a/foreman/test/test_ros_set_goal_server.py b/foreman/test/test_ros_set_goal_server.py new file mode 100644 index 0000000..2acb8b8 --- /dev/null +++ b/foreman/test/test_ros_set_goal_server.py @@ -0,0 +1,115 @@ +import unittest +from unittest.mock import MagicMock + +import rclpy + +from foreman.adapters.ros_set_goal_server import RosSetGoalServer +from foreman.types import ErrorSnapshot +from foreman.types import ForemanErrorCategory +from foreman.types import ForemanResponse +from foreman.types import ForemanSnapshot +from foreman_msgs.srv import SetGoal + + +def _snapshot(goal="force_ctrl", ready=True, at_goal=False, error=None): + """Build a ForemanSnapshot with a no-error default.""" + if error is None: + error = ErrorSnapshot( + is_error=False, + category=ForemanErrorCategory.NONE.value, + message="", + components=[] + ) + return ForemanSnapshot(goal=goal, ready=ready, at_goal=at_goal, error=error, components=[]) + + +def _error_snapshot(message="boom"): + return ErrorSnapshot( + is_error=True, + category=ForemanErrorCategory.EXECUTION.value, + message=message, + components=["ctrl_a"] + ) + + +class TestRosSetGoalServer(unittest.TestCase): + @classmethod + def setUpClass(cls): + rclpy.init() + + @classmethod + def tearDownClass(cls): + rclpy.shutdown() + + def setUp(self): + self.node = rclpy.create_node("test_set_goal_service") + self.addCleanup(self.node.destroy_node) + self.engine = MagicMock() + + def _server(self): + server = RosSetGoalServer(self.node, self.engine) + server._poll_period = 0.0 + return server + + def test_service_is_advertised_as_foreman_set_goal(self): + self._server() + advertised = dict(self.node.get_service_names_and_types()) + self.assertIn("/foreman/set_goal", advertised) + self.assertEqual(advertised["/foreman/set_goal"], ["foreman_msgs/srv/SetGoal"]) + + def test_rejected_goal_returns_without_waiting(self): + self.engine.request_goal.return_value = ForemanResponse( + False, "Goal 'nope' not found in configuration.") + request = SetGoal.Request(goal="nope") + response = SetGoal.Response() + + response = self._server()._handle_set_goal(request, response) + + self.assertFalse(response.success) + self.assertIn("not found", response.message) + self.engine.get_engine_snapshot.assert_not_called() + + def test_succeeds_only_once_engine_reports_at_goal(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.side_effect = [ + _snapshot(at_goal=False), + _snapshot(at_goal=False), + _snapshot(at_goal=True), + ] + request = SetGoal.Request(goal="force_ctrl") + response = SetGoal.Response() + + response = self._server()._handle_set_goal(request, response) + + self.assertTrue(response.success) + self.assertEqual(response.message, "Goal 'force_ctrl' reached.") + self.assertEqual(self.engine.get_engine_snapshot.call_count, 3) + + def test_engine_error_returns_failure(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.side_effect = [ + _snapshot(at_goal=False), + _snapshot(error=_error_snapshot(message="Service rejected the transition.")), + ] + request = SetGoal.Request(goal="force_ctrl") + response = SetGoal.Response() + + response = self._server()._handle_set_goal(request, response) + + self.assertFalse(response.success) + self.assertIn("Service rejected the transition.", response.message) + + def test_preempted_goal_returns_failure(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.return_value = _snapshot(goal="other_goal") + request = SetGoal.Request(goal="force_ctrl") + response = SetGoal.Response() + + response = self._server()._handle_set_goal(request, response) + + self.assertFalse(response.success) + self.assertIn("preempted", response.message) + + +if __name__ == "__main__": + unittest.main() diff --git a/foreman_msgs/CMakeLists.txt b/foreman_msgs/CMakeLists.txt index 57495a3..9580053 100644 --- a/foreman_msgs/CMakeLists.txt +++ b/foreman_msgs/CMakeLists.txt @@ -18,12 +18,17 @@ set(msg_files "msg/ForemanErrorState.msg" ) +set(srv_files + "srv/SetGoal.srv" +) + set(action_files "action/SetGoal.action" ) rosidl_generate_interfaces(${PROJECT_NAME} ${msg_files} + ${srv_files} ${action_files} DEPENDENCIES action_msgs ) diff --git a/foreman_msgs/srv/SetGoal.srv b/foreman_msgs/srv/SetGoal.srv new file mode 100644 index 0000000..fb05de4 --- /dev/null +++ b/foreman_msgs/srv/SetGoal.srv @@ -0,0 +1,6 @@ +# TODO: add description + +string goal +--- +bool success +string message From a1ee86790b82cec653afbb600deb91cc2cebe213 Mon Sep 17 00:00:00 2001 From: Gabriela Vitez Date: Wed, 12 Aug 2026 11:56:07 +0200 Subject: [PATCH 3/4] add shared set_goal lock so actoin and service cannot run goal execution at same time and add reentrant service callback --- .../adapters/ros_set_goal_action_server.py | 101 +++++++++++------- .../foreman/adapters/ros_set_goal_server.py | 85 ++++++++------- foreman/foreman/node.py | 7 +- .../test/test_ros_set_goal_action_server.py | 38 ++++++- foreman/test/test_ros_set_goal_server.py | 43 +++++++- 5 files changed, 189 insertions(+), 85 deletions(-) diff --git a/foreman/foreman/adapters/ros_set_goal_action_server.py b/foreman/foreman/adapters/ros_set_goal_action_server.py index 1a5a367..da91a87 100644 --- a/foreman/foreman/adapters/ros_set_goal_action_server.py +++ b/foreman/foreman/adapters/ros_set_goal_action_server.py @@ -24,10 +24,18 @@ def _to_error_msg(snapshot: ErrorSnapshot) -> ForemanErrorState: class RosSetGoalActionServer: """Drive the system to a named goal with a ROS 2 action.""" - def __init__(self, node: Node, engine: ForemanEngine, poll_period: float = 0.05): + def __init__( + self, + node: Node, + engine: ForemanEngine, + poll_period: float = 0.05, + *, + execution_lock + ): self._node = node self._engine = engine self._poll_period = poll_period + self._execution_lock = execution_lock self.logger_prefix = "Adapters.RosSetGoalActionServer:" self._action_server = ActionServer( @@ -56,53 +64,64 @@ def _execute(self, goal_handle): goal_name = goal_handle.request.goal result = SetGoal.Result() - engine_response = self._engine.request_goal(goal_name) - if not engine_response.success: - self._node.get_logger().warning(f"{engine_response.message}") + if not self._execution_lock.acquire(blocking=False): result.success = False - result.message = engine_response.message + result.message = "Another set_goal request is already active." result.error = _to_error_msg(self._engine.get_engine_snapshot().error) + self._node.get_logger().warning(f"{self.logger_prefix} {result.message}") goal_handle.abort() return result - self._node.get_logger().info(f"{engine_response.message}") - - feedback = SetGoal.Feedback() - while True: - if not goal_handle.is_active: - result.success = False - result.message = f"Goal '{goal_name}' was preempted." - return result - - if goal_handle.is_cancel_requested: + try: + engine_response = self._engine.request_goal(goal_name) + if not engine_response.success: + self._node.get_logger().warning(f"{engine_response.message}") result.success = False - result.message = f"Stopped waiting for goal '{goal_name}'." + result.message = engine_response.message result.error = _to_error_msg(self._engine.get_engine_snapshot().error) - goal_handle.canceled() - return result - - snapshot = self._engine.get_engine_snapshot() - error_msg = _to_error_msg(snapshot.error) - - if snapshot.error.is_error: - result.success = False - result.message = f"[{snapshot.error.category}] {snapshot.error.message}" - result.error = error_msg - self._node.get_logger().error( - f"{self.logger_prefix} Goal '{goal_name}' aborted: {result.message}") goal_handle.abort() return result - if snapshot.at_goal: - result.success = True - result.message = f"Goal '{goal_name}' reached." - result.error = error_msg - self._node.get_logger().info(f"{self.logger_prefix} {result.message}") - goal_handle.succeed() - return result - - feedback.at_goal = False - feedback.error = error_msg - goal_handle.publish_feedback(feedback) - - time.sleep(self._poll_period) + self._node.get_logger().info(f"{engine_response.message}") + + feedback = SetGoal.Feedback() + while True: + if not goal_handle.is_active: + result.success = False + result.message = f"Goal '{goal_name}' was preempted." + return result + + if goal_handle.is_cancel_requested: + result.success = False + result.message = f"Stopped waiting for goal '{goal_name}'." + result.error = _to_error_msg(self._engine.get_engine_snapshot().error) + goal_handle.canceled() + return result + + snapshot = self._engine.get_engine_snapshot() + error_msg = _to_error_msg(snapshot.error) + + if snapshot.error.is_error: + result.success = False + result.message = f"[{snapshot.error.category}] {snapshot.error.message}" + result.error = error_msg + self._node.get_logger().error( + f"{self.logger_prefix} Goal '{goal_name}' aborted: {result.message}") + goal_handle.abort() + return result + + if snapshot.at_goal: + result.success = True + result.message = f"Goal '{goal_name}' reached." + result.error = error_msg + self._node.get_logger().info(f"{self.logger_prefix} {result.message}") + goal_handle.succeed() + return result + + feedback.at_goal = False + feedback.error = error_msg + goal_handle.publish_feedback(feedback) + + time.sleep(self._poll_period) + finally: + self._execution_lock.release() diff --git a/foreman/foreman/adapters/ros_set_goal_server.py b/foreman/foreman/adapters/ros_set_goal_server.py index bc6f6f7..ec22a5f 100644 --- a/foreman/foreman/adapters/ros_set_goal_server.py +++ b/foreman/foreman/adapters/ros_set_goal_server.py @@ -1,6 +1,6 @@ import time -from rclpy.callback_groups import MutuallyExclusiveCallbackGroup +from rclpy.callback_groups import ReentrantCallbackGroup from rclpy.node import Node from foreman.engine import ForemanEngine @@ -10,14 +10,14 @@ class RosSetGoalServer: """ROS 2 service to set a named goal for Foreman Engine.""" - def __init__(self, node: Node, engine: ForemanEngine): + def __init__(self, node: Node, engine: ForemanEngine, *, execution_lock): self._node = node self._engine = engine self._poll_period = 0.05 + self._execution_lock = execution_lock self.logger_prefix = "Adapters.RosSetGoalServer:" - # Using MutuallyExclusiveCallbackGroup - # If a service is processing, we reject new service requests. - self._callback_group = MutuallyExclusiveCallbackGroup() + # Let concurrent callers reach the execution lock and get rejected + self._callback_group = ReentrantCallbackGroup() self._srv = self._node.create_service( SetGoal, @@ -37,42 +37,51 @@ def _handle_set_goal(self, request, response): self._node.get_logger().info( f"{self.logger_prefix} Received request for goal '{goal_name}'") - engine_response = self._engine.request_goal(goal_name) - if not engine_response.success: - self._node.get_logger().warning(f"{engine_response.message}") + if not self._execution_lock.acquire(blocking=False): response.success = False - response.message = engine_response.message + response.message = "Another set_goal request is already active." + self._node.get_logger().warning(f"{self.logger_prefix} {response.message}") return response - self._node.get_logger().info(f"{engine_response.message}") - - while True: - snapshot = self._engine.get_engine_snapshot() - - if snapshot.error.is_error: + try: + engine_response = self._engine.request_goal(goal_name) + if not engine_response.success: + self._node.get_logger().warning(f"{engine_response.message}") response.success = False - response.message = ( - f"[{snapshot.error.category}] {snapshot.error.message}" - ) - self._node.get_logger().error( - f"{self.logger_prefix} Goal '{goal_name}' aborted: " - f"{response.message}") - return response - - if snapshot.goal != goal_name: - response.success = False - response.message = ( - f"Goal '{goal_name}' was preempted by goal '{snapshot.goal}'." - ) - self._node.get_logger().warning( - f"{self.logger_prefix} {response.message}") - return response - - if snapshot.at_goal: - response.success = True - response.message = f"Goal '{goal_name}' reached." - self._node.get_logger().info( - f"{self.logger_prefix} {response.message}") + response.message = engine_response.message return response - time.sleep(self._poll_period) + self._node.get_logger().info(f"{engine_response.message}") + + while True: + snapshot = self._engine.get_engine_snapshot() + + if snapshot.error.is_error: + response.success = False + response.message = ( + f"[{snapshot.error.category}] {snapshot.error.message}" + ) + self._node.get_logger().error( + f"{self.logger_prefix} Goal '{goal_name}' aborted: " + f"{response.message}") + return response + + if snapshot.goal != goal_name: + response.success = False + response.message = ( + f"Goal '{goal_name}' was preempted by goal '{snapshot.goal}'." + ) + self._node.get_logger().warning( + f"{self.logger_prefix} {response.message}") + return response + + if snapshot.at_goal: + response.success = True + response.message = f"Goal '{goal_name}' reached." + self._node.get_logger().info( + f"{self.logger_prefix} {response.message}") + return response + + time.sleep(self._poll_period) + finally: + self._execution_lock.release() diff --git a/foreman/foreman/node.py b/foreman/foreman/node.py index f0a89ac..45d1166 100644 --- a/foreman/foreman/node.py +++ b/foreman/foreman/node.py @@ -22,6 +22,7 @@ def __init__(self): super().__init__('foreman_node') self.foreman_state_lock = threading.Lock() + self.set_goal_execution_lock = threading.Lock() # for error handling ,so we know what and when failed and who to blame self._service_call_active_future = False self._active_transition = None @@ -58,11 +59,13 @@ def __init__(self): self.ros_set_goal_action_server = adapters.RosSetGoalActionServer( node=self, - engine=self.foreman_engine + engine=self.foreman_engine, + execution_lock=self.set_goal_execution_lock ) self.ros_set_goal_server = adapters.RosSetGoalServer( node=self, - engine=self.foreman_engine + engine=self.foreman_engine, + execution_lock=self.set_goal_execution_lock ) # MAIN LOOP ================================================ diff --git a/foreman/test/test_ros_set_goal_action_server.py b/foreman/test/test_ros_set_goal_action_server.py index 7d10592..3d521be 100644 --- a/foreman/test/test_ros_set_goal_action_server.py +++ b/foreman/test/test_ros_set_goal_action_server.py @@ -1,3 +1,4 @@ +import threading import unittest from unittest.mock import MagicMock @@ -73,9 +74,16 @@ def setUp(self): self.addCleanup(self.node.destroy_node) self.engine = MagicMock() - def _server(self): + def _server(self, execution_lock=None): # poll_period=0 so the wait loop does not slow the tests down - return RosSetGoalActionServer(self.node, self.engine, poll_period=0.0) + if execution_lock is None: + execution_lock = threading.Lock() + return RosSetGoalActionServer( + self.node, + self.engine, + poll_period=0.0, + execution_lock=execution_lock + ) def test_action_is_advertised_as_foreman_set_goal(self): self._server() @@ -96,6 +104,20 @@ def test_rejected_goal_aborts_without_waiting(self): self.assertFalse(result.success) self.assertIn("not found", result.message) + def test_busy_set_goal_execution_aborts_without_requesting_goal(self): + self.engine.get_engine_snapshot.return_value = _snapshot() + execution_lock = threading.Lock() + execution_lock.acquire() + handle = _goal_handle() + + result = self._server(execution_lock=execution_lock)._execute(handle) + + handle.abort.assert_called_once() + self.engine.request_goal.assert_not_called() + self.assertFalse(result.success) + self.assertIn("already active", result.message) + execution_lock.release() + def test_succeeds_only_once_engine_reports_at_goal(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") # two polls in transition, then arrived @@ -114,6 +136,18 @@ def test_succeeds_only_once_engine_reports_at_goal(self): # feedback published for each poll that was still in transition self.assertEqual(handle.publish_feedback.call_count, 2) + def test_successful_goal_releases_execution_lock(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=True) + execution_lock = threading.Lock() + handle = _goal_handle() + + result = self._server(execution_lock=execution_lock)._execute(handle) + + self.assertTrue(result.success) + self.assertTrue(execution_lock.acquire(blocking=False)) + execution_lock.release() + def test_feedback_carries_current_error_state(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") self.engine.get_engine_snapshot.side_effect = [ diff --git a/foreman/test/test_ros_set_goal_server.py b/foreman/test/test_ros_set_goal_server.py index 2acb8b8..4c2c384 100644 --- a/foreman/test/test_ros_set_goal_server.py +++ b/foreman/test/test_ros_set_goal_server.py @@ -1,7 +1,9 @@ +import threading import unittest from unittest.mock import MagicMock import rclpy +from rclpy.callback_groups import ReentrantCallbackGroup from foreman.adapters.ros_set_goal_server import RosSetGoalServer from foreman.types import ErrorSnapshot @@ -46,8 +48,14 @@ def setUp(self): self.addCleanup(self.node.destroy_node) self.engine = MagicMock() - def _server(self): - server = RosSetGoalServer(self.node, self.engine) + def _server(self, execution_lock=None): + if execution_lock is None: + execution_lock = threading.Lock() + server = RosSetGoalServer( + self.node, + self.engine, + execution_lock=execution_lock + ) server._poll_period = 0.0 return server @@ -57,6 +65,9 @@ def test_service_is_advertised_as_foreman_set_goal(self): self.assertIn("/foreman/set_goal", advertised) self.assertEqual(advertised["/foreman/set_goal"], ["foreman_msgs/srv/SetGoal"]) + def test_service_uses_reentrant_callback_group(self): + self.assertIsInstance(self._server()._callback_group, ReentrantCallbackGroup) + def test_rejected_goal_returns_without_waiting(self): self.engine.request_goal.return_value = ForemanResponse( False, "Goal 'nope' not found in configuration.") @@ -69,6 +80,20 @@ def test_rejected_goal_returns_without_waiting(self): self.assertIn("not found", response.message) self.engine.get_engine_snapshot.assert_not_called() + def test_busy_set_goal_execution_returns_without_requesting_goal(self): + execution_lock = threading.Lock() + execution_lock.acquire() + request = SetGoal.Request(goal="force_ctrl") + response = SetGoal.Response() + + response = self._server(execution_lock=execution_lock)._handle_set_goal( + request, response) + + self.assertFalse(response.success) + self.assertIn("already active", response.message) + self.engine.request_goal.assert_not_called() + execution_lock.release() + def test_succeeds_only_once_engine_reports_at_goal(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") self.engine.get_engine_snapshot.side_effect = [ @@ -85,6 +110,20 @@ def test_succeeds_only_once_engine_reports_at_goal(self): self.assertEqual(response.message, "Goal 'force_ctrl' reached.") self.assertEqual(self.engine.get_engine_snapshot.call_count, 3) + def test_successful_goal_releases_execution_lock(self): + self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") + self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=True) + execution_lock = threading.Lock() + request = SetGoal.Request(goal="force_ctrl") + response = SetGoal.Response() + + response = self._server(execution_lock=execution_lock)._handle_set_goal( + request, response) + + self.assertTrue(response.success) + self.assertTrue(execution_lock.acquire(blocking=False)) + execution_lock.release() + def test_engine_error_returns_failure(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") self.engine.get_engine_snapshot.side_effect = [ From 3256839ff394fa68087ea3a18bc052b9c1dd3ee9 Mon Sep 17 00:00:00 2001 From: Gabriela Vitez Date: Fri, 14 Aug 2026 12:07:22 +0200 Subject: [PATCH 4/4] Apply reviewer changes --- .../adapters/ros_set_goal_action_server.py | 62 ++++++++----- .../test/test_ros_set_goal_action_server.py | 93 +++++++++++++++++-- foreman_msgs/CMakeLists.txt | 1 + foreman_msgs/action/SetGoal.action | 6 +- foreman_msgs/msg/ComponentState.msg | 8 ++ foreman_msgs/msg/ForemanErrorState.msg | 2 +- 6 files changed, 136 insertions(+), 36 deletions(-) create mode 100644 foreman_msgs/msg/ComponentState.msg diff --git a/foreman/foreman/adapters/ros_set_goal_action_server.py b/foreman/foreman/adapters/ros_set_goal_action_server.py index da91a87..04da09c 100644 --- a/foreman/foreman/adapters/ros_set_goal_action_server.py +++ b/foreman/foreman/adapters/ros_set_goal_action_server.py @@ -6,23 +6,41 @@ from rclpy.node import Node from foreman.engine import ForemanEngine -from foreman.types import ErrorSnapshot +from foreman.types import ForemanSnapshot from foreman_msgs.action import SetGoal +from foreman_msgs.msg import ComponentState from foreman_msgs.msg import ForemanErrorState -def _to_error_msg(snapshot: ErrorSnapshot) -> ForemanErrorState: - """Convert the engine's error snapshot into its ROS representation.""" +def _to_error_msg(snapshot: ForemanSnapshot) -> ForemanErrorState: + """ + Convert the engine's error into its ROS representation. + + The error names the blamed components; their observed states come from the + same snapshot, so a client sees what state each one was in. A component that + is no longer observed is still named, with its state left empty. + """ + observed = {component.name: component for component in snapshot.components} + msg = ForemanErrorState() - msg.is_error = snapshot.is_error - msg.category = snapshot.category - msg.message = snapshot.message - msg.components = list(snapshot.components or []) + msg.is_error = snapshot.error.is_error + msg.category = snapshot.error.category + msg.message = snapshot.error.message + + for name in snapshot.error.components or []: + component_msg = ComponentState() + component_msg.name = name + component = observed.get(name) + if component: + component_msg.component_type = component.component_type.value + component_msg.lifecycle_state = component.lifecycle_state.name + msg.components.append(component_msg) + return msg class RosSetGoalActionServer: - """Drive the system to a named goal with a ROS 2 action.""" + """ROS 2 action interface to set the Foreman goal.""" def __init__( self, @@ -32,11 +50,10 @@ def __init__( *, execution_lock ): - self._node = node self._engine = engine self._poll_period = poll_period self._execution_lock = execution_lock - self.logger_prefix = "Adapters.RosSetGoalActionServer:" + self._logger = node.get_logger().get_child('action') self._action_server = ActionServer( node, @@ -48,11 +65,10 @@ def __init__( callback_group=node.callback_group_subscriber ) - self._node.get_logger().info(f"{self.logger_prefix} Action /foreman/set_goal is ready.") + self._logger.info("Action /foreman/set_goal is ready.") def _on_goal_request(self, goal_request) -> GoalResponse: - self._node.get_logger().info( - f"{self.logger_prefix} Received request for goal '{goal_request.goal}'") + self._logger.debug(f"Received request for goal '{goal_request.goal}'") return GoalResponse.ACCEPT def _on_cancel_request(self, goal_handle) -> CancelResponse: @@ -67,46 +83,44 @@ def _execute(self, goal_handle): if not self._execution_lock.acquire(blocking=False): result.success = False result.message = "Another set_goal request is already active." - result.error = _to_error_msg(self._engine.get_engine_snapshot().error) - self._node.get_logger().warning(f"{self.logger_prefix} {result.message}") + self._logger.warning(result.message) goal_handle.abort() return result try: engine_response = self._engine.request_goal(goal_name) if not engine_response.success: - self._node.get_logger().warning(f"{engine_response.message}") + self._logger.warning(engine_response.message) result.success = False result.message = engine_response.message - result.error = _to_error_msg(self._engine.get_engine_snapshot().error) goal_handle.abort() return result - self._node.get_logger().info(f"{engine_response.message}") + self._logger.debug(engine_response.message) feedback = SetGoal.Feedback() while True: if not goal_handle.is_active: + # Reachable when the action server is destroyed on shutdown: result.success = False - result.message = f"Goal '{goal_name}' was preempted." + result.message = f"Goal '{goal_name}' is no longer active." return result if goal_handle.is_cancel_requested: result.success = False result.message = f"Stopped waiting for goal '{goal_name}'." - result.error = _to_error_msg(self._engine.get_engine_snapshot().error) + result.error = _to_error_msg(self._engine.get_engine_snapshot()) goal_handle.canceled() return result snapshot = self._engine.get_engine_snapshot() - error_msg = _to_error_msg(snapshot.error) + error_msg = _to_error_msg(snapshot) if snapshot.error.is_error: result.success = False result.message = f"[{snapshot.error.category}] {snapshot.error.message}" result.error = error_msg - self._node.get_logger().error( - f"{self.logger_prefix} Goal '{goal_name}' aborted: {result.message}") + self._logger.error(f"Goal '{goal_name}' aborted: {result.message}") goal_handle.abort() return result @@ -114,7 +128,7 @@ def _execute(self, goal_handle): result.success = True result.message = f"Goal '{goal_name}' reached." result.error = error_msg - self._node.get_logger().info(f"{self.logger_prefix} {result.message}") + self._logger.info(result.message) goal_handle.succeed() return result diff --git a/foreman/test/test_ros_set_goal_action_server.py b/foreman/test/test_ros_set_goal_action_server.py index 3d521be..15ef4fa 100644 --- a/foreman/test/test_ros_set_goal_action_server.py +++ b/foreman/test/test_ros_set_goal_action_server.py @@ -8,13 +8,24 @@ from foreman.adapters.ros_set_goal_action_server import _to_error_msg from foreman.adapters.ros_set_goal_action_server import RosSetGoalActionServer +from foreman.types import Component +from foreman.types import ComponentType from foreman.types import ErrorSnapshot from foreman.types import ForemanErrorCategory from foreman.types import ForemanResponse from foreman.types import ForemanSnapshot +from foreman.types import LifecycleState -def _snapshot(goal="force_ctrl", ready=True, at_goal=False, error=None): +def _component(name="ctrl_a", state=LifecycleState.INACTIVE): + return Component( + name=name, + component_type=ComponentType.CONTROLLER, + lifecycle_state=state + ) + + +def _snapshot(goal="force_ctrl", ready=True, at_goal=False, error=None, components=None): """Build a ForemanSnapshot with a no-error default.""" if error is None: error = ErrorSnapshot( @@ -23,7 +34,13 @@ def _snapshot(goal="force_ctrl", ready=True, at_goal=False, error=None): message="", components=[] ) - return ForemanSnapshot(goal=goal, ready=ready, at_goal=at_goal, error=error, components=[]) + return ForemanSnapshot( + goal=goal, + ready=ready, + at_goal=at_goal, + error=error, + components=components if components is not None else [] + ) def _error_snapshot(category=ForemanErrorCategory.EXECUTION, message="boom", components=None): @@ -46,17 +63,41 @@ def _goal_handle(goal_name="force_ctrl"): class TestToErrorMsg(unittest.TestCase): def test_maps_every_field(self): - msg = _to_error_msg(_error_snapshot(components=["a", "b"])) + msg = _to_error_msg(_snapshot(error=_error_snapshot(components=["a", "b"]))) self.assertTrue(msg.is_error) self.assertEqual(msg.category, ForemanErrorCategory.EXECUTION.value) self.assertEqual(msg.message, "boom") - self.assertEqual(list(msg.components), ["a", "b"]) + self.assertEqual([c.name for c in msg.components], ["a", "b"]) def test_none_components_become_empty_list(self): - msg = _to_error_msg(ErrorSnapshot( - is_error=False, category="None", message="", components=None)) + msg = _to_error_msg(_snapshot(error=ErrorSnapshot( + is_error=False, category="None", message="", components=None))) self.assertEqual(list(msg.components), []) + def test_blamed_components_carry_their_observed_state(self): + msg = _to_error_msg(_snapshot( + error=_error_snapshot(components=["ctrl_a"]), + components=[_component("ctrl_a", LifecycleState.ACTIVE)] + )) + + self.assertEqual(len(msg.components), 1) + self.assertEqual(msg.components[0].name, "ctrl_a") + self.assertEqual(msg.components[0].component_type, ComponentType.CONTROLLER.value) + self.assertEqual(msg.components[0].lifecycle_state, "ACTIVE") + + def test_vanished_component_is_named_without_a_state(self): + # "Required components vanished from /activity" drops the component from + # the observed state, so only its name is known. + msg = _to_error_msg(_snapshot( + error=_error_snapshot(components=["gone"]), + components=[_component("still_here")] + )) + + self.assertEqual(len(msg.components), 1) + self.assertEqual(msg.components[0].name, "gone") + self.assertEqual(msg.components[0].component_type, "") + self.assertEqual(msg.components[0].lifecycle_state, "") + class TestRosSetGoalActionServer(unittest.TestCase): @classmethod @@ -104,6 +145,20 @@ def test_rejected_goal_aborts_without_waiting(self): self.assertFalse(result.success) self.assertIn("not found", result.message) + def test_rejected_goal_does_not_report_a_leftover_engine_error(self): + self.engine.request_goal.return_value = ForemanResponse( + False, "Goal 'nope' not found in configuration.") + self.engine.get_engine_snapshot.return_value = _snapshot( + error=_error_snapshot(message="leftover failure") + ) + + result = self._server()._execute(_goal_handle("nope")) + + self.assertFalse(result.success) + self.assertIn("not found", result.message) + self.assertFalse(result.error.is_error) + self.assertEqual(result.error.message, "") + def test_busy_set_goal_execution_aborts_without_requesting_goal(self): self.engine.get_engine_snapshot.return_value = _snapshot() execution_lock = threading.Lock() @@ -118,6 +173,21 @@ def test_busy_set_goal_execution_aborts_without_requesting_goal(self): self.assertIn("already active", result.message) execution_lock.release() + def test_busy_set_goal_execution_does_not_report_an_engine_error(self): + self.engine.get_engine_snapshot.return_value = _snapshot( + error=_error_snapshot(message="unrelated failure") + ) + execution_lock = threading.Lock() + execution_lock.acquire() + + result = self._server(execution_lock=execution_lock)._execute(_goal_handle()) + + self.assertFalse(result.success) + self.assertIn("already active", result.message) + self.assertFalse(result.error.is_error) + self.assertEqual(result.error.message, "") + execution_lock.release() + def test_succeeds_only_once_engine_reports_at_goal(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") # two polls in transition, then arrived @@ -166,7 +236,10 @@ def test_engine_error_aborts_and_reports_it(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") self.engine.get_engine_snapshot.side_effect = [ _snapshot(at_goal=False), - _snapshot(error=_error_snapshot(message="Service rejected the transition.")), + _snapshot( + error=_error_snapshot(message="Service rejected the transition."), + components=[_component("ctrl_a", LifecycleState.INACTIVE)] + ), ] handle = _goal_handle() @@ -178,7 +251,8 @@ def test_engine_error_aborts_and_reports_it(self): self.assertIn("Service rejected the transition.", result.message) self.assertTrue(result.error.is_error) self.assertEqual(result.error.category, ForemanErrorCategory.EXECUTION.value) - self.assertEqual(list(result.error.components), ["ctrl_a"]) + self.assertEqual([c.name for c in result.error.components], ["ctrl_a"]) + self.assertEqual(result.error.components[0].lifecycle_state, "INACTIVE") def test_cancel_stops_waiting(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") @@ -193,7 +267,7 @@ def test_cancel_stops_waiting(self): handle.abort.assert_not_called() self.assertFalse(result.success) - def test_preempted_goal_returns_without_terminal_call(self): + def test_inactive_goal_returns_without_terminal_call(self): self.engine.request_goal.return_value = ForemanResponse(True, "Goal accepted.") self.engine.get_engine_snapshot.return_value = _snapshot(at_goal=False) handle = _goal_handle() @@ -205,6 +279,7 @@ def test_preempted_goal_returns_without_terminal_call(self): handle.abort.assert_not_called() handle.canceled.assert_not_called() self.assertFalse(result.success) + self.assertIn("no longer active", result.message) def test_already_at_goal_succeeds_without_feedback(self): self.engine.request_goal.return_value = ForemanResponse( diff --git a/foreman_msgs/CMakeLists.txt b/foreman_msgs/CMakeLists.txt index 9580053..0847133 100644 --- a/foreman_msgs/CMakeLists.txt +++ b/foreman_msgs/CMakeLists.txt @@ -15,6 +15,7 @@ find_package(rosidl_default_generators REQUIRED) find_package(action_msgs REQUIRED) set(msg_files + "msg/ComponentState.msg" "msg/ForemanErrorState.msg" ) diff --git a/foreman_msgs/action/SetGoal.action b/foreman_msgs/action/SetGoal.action index 874b5a8..81f7870 100644 --- a/foreman_msgs/action/SetGoal.action +++ b/foreman_msgs/action/SetGoal.action @@ -1,10 +1,12 @@ -# Drive the system to a named goal state. -# Succeeds once the goal is reached, not when it is accepted. +# Set the Foreman goal and wait until the system reaches it. +# Succeeds once the goal is reached string goal --- bool success +# Why the goal ended: reached, refused, preempted, cancelled or aborted. string message +# Engine error state when the goal ended. Empty if lock is busy or engine rejected goal ForemanErrorState error --- bool at_goal diff --git a/foreman_msgs/msg/ComponentState.msg b/foreman_msgs/msg/ComponentState.msg new file mode 100644 index 0000000..1345699 --- /dev/null +++ b/foreman_msgs/msg/ComponentState.msg @@ -0,0 +1,8 @@ +# A component and the state it was observed in. + +string name +# ComponentType value: hardware, controller or lifecycle_node. +# Empty when the component is no longer observed, e.g. it vanished from +# /activity, so the name is still reported without a known state. +string component_type +string lifecycle_state diff --git a/foreman_msgs/msg/ForemanErrorState.msg b/foreman_msgs/msg/ForemanErrorState.msg index ccceab7..beb1bb3 100644 --- a/foreman_msgs/msg/ForemanErrorState.msg +++ b/foreman_msgs/msg/ForemanErrorState.msg @@ -6,4 +6,4 @@ bool is_error # UnexpectedStateError, PlannerError, or None. string category string message -string[] components +ComponentState[] components