Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 28 additions & 12 deletions .pre-commit-config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -2,21 +2,37 @@ repos:
- repo: https://github.com/pre-commit/pre-commit-hooks
rev: v6.0.0
hooks:
- id: trailing-whitespace
- id: end-of-file-fixer
- id: check-added-large-files
- id: check-ast
- id: check-case-conflict
- id: check-docstring-first
- id: check-merge-conflict
- id: check-symlinks
- id: check-xml
- id: check-yaml
- repo: https://github.com/PyCQA/isort
rev: 8.0.1
hooks:
- id: isort
args: ["--profile", "google", "--line-length=99"]
- repo: https://github.com/hhatto/autopep8
rev: v2.3.2
args: ["--allow-multiple-documents"]
- id: debug-statements
- id: end-of-file-fixer
- id: mixed-line-ending
- id: trailing-whitespace
- id: fix-byte-order-marker
- repo: https://github.com/asottile/pyupgrade
rev: v3.21.2
hooks:
- id: autopep8
args: ["-i", "--max-line-length=99"]
- id: pyupgrade
args: [--py36-plus]
- repo: https://github.com/PyCQA/pydocstyle
rev: 6.3.0
hooks:
- id: pydocstyle
args: ["--convention=pep257", "--add-ignore=D100,D101,D102,D103,D104,D105,D106,D107"]
args: ["--ignore=D100,D101,D102,D103,D104,D105,D106,D107,D203,D212,D404"]
- repo: https://github.com/psf/black
rev: 26.5.1
hooks:
- id: black
args: ["--line-length=99"]
- repo: https://github.com/pycqa/flake8
rev: 7.3.0
hooks:
- id: flake8
args: ["--extend-ignore=E501"]
29 changes: 16 additions & 13 deletions foreman/foreman/adapters/autostart_adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@
class AutostartAdapter:
"""Adapter to transition automatically to a desired state after all desired components are loaded."""

STABLE_TICKS_REQUIRED = 50 # consecutive ticks with no state change before requesting transition
STABLE_TICKS_REQUIRED = (
50 # consecutive ticks with no state change before requesting transition
)

def __init__(self, node: Node, engine: ForemanEngine, goal_name: str, autostart: bool = False):
self._node = node
Expand All @@ -20,7 +22,7 @@ def __init__(self, node: Node, engine: ForemanEngine, goal_name: str, autostart:

def autostart(self):
if not self._autostart:
self._node.get_logger().info(f"Autostart parameter is set to false.")
self._node.get_logger().info("Autostart parameter is set to false.")
return

if self.transition_success and self.engine.get_engine_snapshot().error.is_error:
Expand All @@ -35,14 +37,18 @@ def autostart(self):
if not self.all_components_ready():
not_ready = self._get_not_ready_components()
self._node.get_logger().info(
f"Some components are not ready yet: {not_ready}. Waiting before requesting autostart...", throttle_duration_sec=2)
f"Some components are not ready yet: {not_ready}. Waiting before requesting autostart...",
throttle_duration_sec=2,
)
self._stable_ticks = 0
return

if not self._is_state_stable():
return

self._node.get_logger().info(f"All components stable and ready. Requesting goal transition.")
self._node.get_logger().info(
"All components stable and ready. Requesting goal transition."
)
self.transition_success = self.send_goal_request()

@property
Expand All @@ -64,7 +70,8 @@ def send_goal_request(self) -> bool:
def _is_state_stable(self) -> bool:
"""Return True once the system state has been unchanged for STABLE_TICKS_REQUIRED consecutive ticks."""
current_states = {
c.name: c.lifecycle_state for c in self.engine.get_engine_snapshot().components}
c.name: c.lifecycle_state for c in self.engine.get_engine_snapshot().components
}
if current_states != self._last_observed_states:
self._last_observed_states = current_states
self._stable_ticks = 0
Expand All @@ -77,19 +84,15 @@ def _get_not_ready_components(self, desired_state=LifecycleState.UNCONFIGURED):
snapshot = self.engine.get_engine_snapshot()
observed = {c.name: c for c in snapshot.components}
missing = [n for n in self.engine._config.tracked_components if n not in observed]
wrong_state = [
c.name for c in observed.values()
if c.lifecycle_state < desired_state
]
wrong_state = [c.name for c in observed.values() if c.lifecycle_state < desired_state]
return missing + wrong_state

def all_components_ready(self, desired_state=LifecycleState.UNCONFIGURED):
"""Check if all the tracked components are in the desired state."""
snapshot = self.engine.get_engine_snapshot()
observed = {c.name: c for c in snapshot.components}

all_unconfigured = (
all(name in observed for name in self.engine._config.tracked_components)
and all(c.lifecycle_state >= desired_state for c in observed.values())
)
all_unconfigured = all(
name in observed for name in self.engine._config.tracked_components
) and all(c.lifecycle_state >= desired_state for c in observed.values())
return all_unconfigured
46 changes: 19 additions & 27 deletions foreman/foreman/adapters/component_state_monitor.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,12 @@
from controller_manager_msgs.msg import ControllerManagerActivity
from lifecycle_msgs.msg import TransitionEvent
from lifecycle_msgs.srv import GetState
from rclpy.event_handler import QoSSubscriptionMatchedInfo
from rclpy.event_handler import SubscriptionEventCallbacks
from rclpy.event_handler import QoSSubscriptionMatchedInfo, SubscriptionEventCallbacks
from rclpy.node import Node
from rclpy.qos import DurabilityPolicy
from rclpy.qos import HistoryPolicy
from rclpy.qos import QoSProfile
from rclpy.qos import ReliabilityPolicy
from rclpy.qos import DurabilityPolicy, HistoryPolicy, QoSProfile, ReliabilityPolicy

from foreman.engine import ForemanEngine
from foreman.types import Component
from foreman.types import ComponentType
from foreman.types import LifecycleState
from foreman.types import Component, ComponentType, LifecycleState


class ComponentStateMonitor:
Expand Down Expand Up @@ -51,14 +45,14 @@ def __init__(
reliability=ReliabilityPolicy.RELIABLE,
durability=DurabilityPolicy.TRANSIENT_LOCAL,
history=HistoryPolicy.KEEP_LAST,
depth=1
depth=1,
)
self._subscription = self._node.create_subscription(
ControllerManagerActivity,
f'/{controller_manager_name}/activity',
f"/{controller_manager_name}/activity",
self._activity_callback,
qos_profile,
callback_group=self._node.callback_group_subscriber
callback_group=self._node.callback_group_subscriber,
)
self._node.get_logger().info(
f"{self._logger_prefix} Subscribed to /{controller_manager_name}/activity"
Expand All @@ -70,8 +64,8 @@ def __init__(
for lc_node_name in lifecycle_nodes:
self._lc_node_get_state_clients[lc_node_name] = self._node.create_client(
GetState,
f'/{lc_node_name}/get_state',
callback_group=self._node.callback_group_services
f"/{lc_node_name}/get_state",
callback_group=self._node.callback_group_services,
)

# this matched event on transition_event/ topic handles both
Expand All @@ -83,12 +77,13 @@ def __init__(
)
self._node.create_subscription(
TransitionEvent,
f'/{lc_node_name}/transition_event',
f"/{lc_node_name}/transition_event",
callback=lambda msg, n=lc_node_name: self._lifecycle_transition_event_callback(
n, msg),
n, msg
),
qos_profile=10,
event_callbacks=event_callbacks,
callback_group=self._node.callback_group_subscriber
callback_group=self._node.callback_group_subscriber,
)

if lifecycle_nodes:
Expand All @@ -107,7 +102,7 @@ def _activity_callback(self, msg: ControllerManagerActivity):
components[hw_msg.name] = Component(
name=hw_msg.name,
component_type=ComponentType.HARDWARE,
lifecycle_state=LifecycleState(hw_msg.state.id)
lifecycle_state=LifecycleState(hw_msg.state.id),
)
except ValueError:
continue
Expand All @@ -117,7 +112,7 @@ def _activity_callback(self, msg: ControllerManagerActivity):
components[ctrl_msg.name] = Component(
name=ctrl_msg.name,
component_type=ComponentType.CONTROLLER,
lifecycle_state=LifecycleState(ctrl_msg.state.id)
lifecycle_state=LifecycleState(ctrl_msg.state.id),
)
except ValueError:
continue
Expand All @@ -138,7 +133,7 @@ def _on_lifecycle_publisher_matched(self, name: str, info: QoSSubscriptionMatche
self._lc_components[name] = Component(
name=name,
component_type=ComponentType.LIFECYCLE_NODE,
lifecycle_state=LifecycleState.FINALIZED
lifecycle_state=LifecycleState.FINALIZED,
)
self._node.get_logger().warning(
f"{self._logger_prefix} Lifecycle node '{name}' disconnected."
Expand All @@ -157,9 +152,7 @@ def _on_lifecycle_get_state_response(self, name: str, future):
response = future.result()
state = LifecycleState(response.current_state.id)
self._lc_components[name] = Component(
name=name,
component_type=ComponentType.LIFECYCLE_NODE,
lifecycle_state=state
name=name, component_type=ComponentType.LIFECYCLE_NODE, lifecycle_state=state
)
self._node.get_logger().info(
f"{self._logger_prefix} Lifecycle node '{name}' discovered. State: {state.name}"
Expand All @@ -176,9 +169,7 @@ def _lifecycle_transition_event_callback(self, name: str, msg: TransitionEvent):
try:
new_state = LifecycleState(msg.goal_state.id)
self._lc_components[name] = Component(
name=name,
component_type=ComponentType.LIFECYCLE_NODE,
lifecycle_state=new_state
name=name, component_type=ComponentType.LIFECYCLE_NODE, lifecycle_state=new_state
)
self._push_merged_state()
except ValueError:
Expand All @@ -193,7 +184,8 @@ def _push_merged_state(self):

if not was_ready and self._engine.is_ready:
self._node.get_logger().info(
f"{self._logger_prefix} Foreman is READY. Fresh state received.")
f"{self._logger_prefix} Foreman is READY. Fresh state received."
)

if not response.success and response.error:
self._node.get_logger().error(
Expand Down
59 changes: 36 additions & 23 deletions foreman/foreman/adapters/controller_manager_service_caller.py
Original file line number Diff line number Diff line change
@@ -1,18 +1,13 @@
import time
from typing import List

from controller_manager_msgs.srv import CleanupController
from controller_manager_msgs.srv import ConfigureController
from controller_manager_msgs.srv import SetHardwareComponentState
from controller_manager_msgs.srv import SwitchController
from controller_manager_msgs.srv import (
CleanupController,
ConfigureController,
SetHardwareComponentState,
SwitchController,
)
from rclpy.node import Node
from rclpy.task import Future

from foreman.types import ComponentType
from foreman.types import ForemanError
from foreman.types import ForemanErrorCategory
from foreman.types import LifecycleState
from foreman.types import SystemTransitionCommand
from foreman.types import ComponentType, LifecycleState, SystemTransitionCommand


class ControllerManagerServiceCaller:
Expand All @@ -26,21 +21,33 @@ def __init__(self, node: Node, controller_manager_name: str):
group = self._node.callback_group_services

self._client_set_hardware_component_state = self._node.create_client(
SetHardwareComponentState, f'/{controller_manager_name}/set_hardware_component_state', callback_group=group)
SetHardwareComponentState,
f"/{controller_manager_name}/set_hardware_component_state",
callback_group=group,
)
self._client_configure_controller = self._node.create_client(
ConfigureController, f'/{controller_manager_name}/configure_controller', callback_group=group)
ConfigureController,
f"/{controller_manager_name}/configure_controller",
callback_group=group,
)
self._client_cleanup_controller = self._node.create_client(
CleanupController, f'/{controller_manager_name}/cleanup_controller', callback_group=group)
CleanupController,
f"/{controller_manager_name}/cleanup_controller",
callback_group=group,
)
self._client_switch_controller = self._node.create_client(
SwitchController, f'/{controller_manager_name}/switch_controller', callback_group=group)
SwitchController, f"/{controller_manager_name}/switch_controller", callback_group=group
)

self._node.get_logger().info(
f"{self.logger_prefix} {self._controller_manager_name} service clients created.")
f"{self.logger_prefix} {self._controller_manager_name} service clients created."
)

def _service_call(self, client, request) -> Future:
if not client.service_is_ready():
raise RuntimeError(
f"Service {client.srv_name} not ready. Is {self._controller_manager_name} running?")
f"Service {client.srv_name} not ready. Is {self._controller_manager_name} running?"
)
return client.call_async(request)

def execute_transition(self, cmd: SystemTransitionCommand) -> Future:
Expand All @@ -61,21 +68,27 @@ def execute_transition(self, cmd: SystemTransitionCommand) -> Future:
if goal == LifecycleState.ACTIVE:
self._node.get_logger().info(f"{self.logger_prefix} Switch -> Activate {name}")
req = SwitchController.Request(
activate_controllers=[name], strictness=SwitchController.Request.STRICT)
activate_controllers=[name], strictness=SwitchController.Request.STRICT
)
return self._service_call(self._client_switch_controller, req)

elif goal == LifecycleState.INACTIVE and current == LifecycleState.ACTIVE:
self._node.get_logger().info(f"{self.logger_prefix} Switch -> Deactivate {name}")
req = SwitchController.Request(deactivate_controllers=[
name], strictness=SwitchController.Request.STRICT)
req = SwitchController.Request(
deactivate_controllers=[name], strictness=SwitchController.Request.STRICT
)
return self._service_call(self._client_switch_controller, req)

elif goal == LifecycleState.INACTIVE and current == LifecycleState.UNCONFIGURED:
self._node.get_logger().info(f"{self.logger_prefix} Configure -> {name}")
return self._service_call(self._client_configure_controller, ConfigureController.Request(name=name))
return self._service_call(
self._client_configure_controller, ConfigureController.Request(name=name)
)

elif goal == LifecycleState.UNCONFIGURED:
self._node.get_logger().info(f"{self.logger_prefix} Cleanup -> {name}")
return self._service_call(self._client_cleanup_controller, CleanupController.Request(name=name))
return self._service_call(
self._client_cleanup_controller, CleanupController.Request(name=name)
)

raise ValueError(f"Unable to process transition command: {cmd}")
15 changes: 4 additions & 11 deletions foreman/foreman/adapters/lifecycle_node_service_caller.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,7 @@
from rclpy.node import Node
from rclpy.task import Future

from foreman.types import LifecycleState
from foreman.types import SystemTransitionCommand
from foreman.types import LifecycleState, SystemTransitionCommand

# maps (current_state, goal_state) to the lifecycle transition ID
_TRANSITION_MAP = {
Expand All @@ -30,9 +29,7 @@ def __init__(self, node: Node, lifecycle_nodes: List[str]):

for lc_name in lifecycle_nodes:
client = self._node.create_client(
ChangeState,
f'/{lc_name}/change_state',
callback_group=group
ChangeState, f"/{lc_name}/change_state", callback_group=group
)
self._clients[lc_name] = client

Expand All @@ -52,19 +49,15 @@ def execute_transition(self, cmd: SystemTransitionCommand) -> Future:
raise ValueError(f"No lifecycle client for node '{name}'")

if not client.service_is_ready():
raise RuntimeError(
f"Service /{name}/change_state not ready. Is '{name}' running?"
)
raise RuntimeError(f"Service /{name}/change_state not ready. Is '{name}' running?")

transition_id = _TRANSITION_MAP.get((current, goal))
if transition_id is None:
raise ValueError(
f"No valid lifecycle transition from {current.name} to {goal.name} for '{name}'"
)

self._node.get_logger().info(
f"{self.logger_prefix} {name} -> {goal.name}"
)
self._node.get_logger().info(f"{self.logger_prefix} {name} -> {goal.name}")

req = ChangeState.Request()
req.transition.id = transition_id
Expand Down
Loading