Index: leahi_dialin/dd/dialysate_delivery.py =================================================================== diff -u -r79697500614904bafb56f3d33e368b954ba78f7d -r88eea899fa8d03596f944505285bf3049e64312b --- leahi_dialin/dd/dialysate_delivery.py (.../dialysate_delivery.py) (revision 79697500614904bafb56f3d33e368b954ba78f7d) +++ leahi_dialin/dd/dialysate_delivery.py (.../dialysate_delivery.py) (revision 88eea899fa8d03596f944505285bf3049e64312b) @@ -7,20 +7,19 @@ # # @file dialysate_delivery.py # -# @author (last) Dara Navaei -# @date (last) 26-Feb-2024 +# @author (last) Zoltan Miskolci +# @date (last) 04-May-2026 # @author (original) Peter Lucia # @date (original) 02-Apr-2020 # ############################################################################ -import struct +# Project imports from .modules.alarms import DDAlarms from .modules.balancing_chamber import DDBalancingChamber from .modules.blood_leak import DDBloodLeak from .modules.concentrate_pump import DDConcentratePumps from .modules.conductivity_sensors import DDConductivitySensors -from .modules.constants import NO_RESET, RESET from .modules.dialysate_pump import DDDialysatePumps from .modules.events import DDEvents from .modules.gen_dialysate import DDGenDialysate @@ -31,6 +30,9 @@ from .modules.pre_gen_dialysate import DDPreGenDialysate from .modules.rinse_pump import DDRinsePump from .modules.spent_chamber_fill import DDSpentChamberFill +from .modules.substitution_pump import DDSubstitutionPump +from .modules.drybicart import DDDryBicart +from .modules.mixing_cntrl import DDDialysateMixing from .modules.temperature_sensors import DDTemperatureSensors from .modules.dd_test_configs import DDTestConfig from .modules.ultrafiltration import DDUltrafiltration @@ -40,13 +42,15 @@ from .proxies.ro_proxy import ROProxy from .proxies.td_proxy import TDProxy -from ..common.msg_defs import MsgIds, MsgFieldPositions, MsgFieldPositionsFWVersions -from ..common.dd_defs import DDOpModes -from ..protocols.CAN import DenaliMessage, DenaliCanMessenger, DenaliChannels -from ..utils.base import AbstractSubSystem, publish, LogManager -from ..utils.checks import check_broadcast_interval_override_ms -from ..utils.conversions import integer_to_bytearray, unsigned_short_to_bytearray, bytearray_to_integer, \ - bytearray_to_byte +from ..common.constants import NO_RESET +from ..common import dd_enum_repository +from ..common.generic_defs import DataTypes +from ..common.msg_defs import MsgIds, MsgFieldPositions +from ..common.override_templates import cmd_generic_broadcast_interval_override, cmd_generic_override +from ..protocols.CAN import CanMessage, CanMessenger, CanChannels +from leahi_dialin.utils.abstract_classes import AbstractSubSystem +from leahi_dialin.utils.base import publish, LogManager +from ..utils.conversions import integer_to_bytearray, bytearray_to_byte class DD(AbstractSubSystem): @@ -74,76 +78,88 @@ @param can_interface: (str) CANBus interface name, e.g. "can0" @param log_level: (str) Logging level, defaults to None """ - super().__init__() self._log_manager = LogManager(log_level=log_level, log_filepath=self.__class__.__name__ + ".log") self.logger = self._log_manager.logger # Create listener - self.can_interface = DenaliCanMessenger(can_interface=can_interface, + self.can_interface = CanMessenger(can_interface=can_interface, logger=self.logger) self.can_interface.start() self.callback_id = None # register handler for DD operation mode broadcast messages if self.can_interface is not None: - channel_id = DenaliChannels.dd_sync_broadcast_ch_id - self.msg_id_dd_op_mode_data = MsgIds.MSG_ID_DD_OP_MODE_DATA.value - self.can_interface.register_receiving_publication_function(channel_id, self.msg_id_dd_op_mode_data, - self._handler_dd_op_mode_sync) + self.can_interface.register_receiving_publication_function(channel_id = CanChannels.dd_sync_broadcast_ch_id, + message_id = MsgIds.MSG_ID_DD_OP_MODE_DATA.value, + function = self._handler_dd_op_mode_sync) - self.msg_id_dd_version_response = MsgIds.MSG_ID_DD_VERSION_RESPONSE.value - self.can_interface.register_receiving_publication_function(channel_id, - self.msg_id_dd_version_response, - self._handler_dd_version_response_sync) + self.can_interface.register_receiving_publication_function(channel_id = CanChannels.dd_sync_broadcast_ch_id, + message_id = MsgIds.MSG_ID_DD_VERSION_RESPONSE.value, + function = self._handler_dd_version_response_sync) - self.msg_id_dd_debug_event = MsgIds.MSG_ID_DD_DEBUG_EVENT.value - self.can_interface.register_receiving_publication_function(channel_id, - self.msg_id_dd_debug_event, - self._handler_dd_debug_event_sync) + self.can_interface.register_receiving_publication_function(channel_id = CanChannels.dd_sync_broadcast_ch_id, + message_id = MsgIds.MSG_ID_DD_DEBUG_EVENT.value, + function = self._handler_dd_debug_event_sync) + # Dialin will send a login message during construction. This is for the leahi subsystems to start + # publishing CAN data when there is no UI connected as the UI typically does this job. + self.cmd_log_in_to_dd() + # create properties - self.dd_op_mode_timestamp = 0.0 - self.dd_debug_events_timestamp = 0.0 - self.dd_version_response_timestamp = 0.0 - self.dd_operation_mode = DDOpModes.MODE_INIT.value - self.dd_operation_sub_mode = 0 - self.dd_logged_in = False + self.dd_op_mode_timestamp = 0.0 #: The timestamp of the latest operation mode message + self.dd_debug_events_timestamp = 0.0 #: The timestamp of the latest events message + self.dd_version_response_timestamp = 0.0 #: The timestamp of the latest DD version info message + self.dd_operation_mode = dd_enum_repository.DDOpModes.MODE_INIT.value #: The Operation Mode's value + self.dd_operation_sub_mode = 0 #: The Operation Sub-Mode's value + self.dd_logged_in = False #: The value showing if the user is logged in or not self.dd_set_logged_in_status(False) - self.dd_version = None - self.dd_fpga_version = None - self.dd_debug_events = [''] * self._DD_DEBUG_EVENT_LIST_COUNT - self.dd_debug_event_index = 0 - self.dd_last_debug_event = '' + self.dd_version = None #: The DD's version value + self.dd_fpga_version = None #: The DD's FPGA version value + self.dd_debug_events = [''] * self._DD_DEBUG_EVENT_LIST_COUNT #: The Debug Event's list + self.dd_debug_event_index = 0 #: The index of the last Event + self.dd_last_debug_event = '' #: The name of the last Event # Create command groups - self.alarms = DDAlarms(self.can_interface, self.logger) - self.balancing_chamber = DDBalancingChamber(self.can_interface, self.logger) - self.blood_leak = DDBloodLeak(self.can_interface, self.logger) - self.concentrate_pumps = DDConcentratePumps(self.can_interface, self.logger) - self.conductivity_sensors = DDConductivitySensors(self.can_interface, self.logger) - self.dialysate_pumps = DDDialysatePumps(self.can_interface, self.logger) - self.events = DDEvents(self.can_interface, self.logger) - self.gen_dialysate = DDGenDialysate(self.can_interface, self.logger) - self.heaters = DDHeaters(self.can_interface, self.logger) - self.levels = DDLevels(self.can_interface, self.logger) - self.post_gen_dialysate = DDPostGenDialysate(self.can_interface, self.logger) - self.pressure_sensors = DDPressureSensors(self.can_interface, self.logger) - self.pre_gen_dialysate = DDPreGenDialysate(self.can_interface, self.logger) - self.rinse_pump = DDRinsePump(self.can_interface, self.logger) - self.spent_chamber_fill = DDSpentChamberFill(self.can_interface, self.logger) - self.temperature_sensors = DDTemperatureSensors(self.can_interface, self.logger) - self.test_configs = DDTestConfig(self.can_interface, self.logger) - self.ultrafiltration = DDUltrafiltration(self.can_interface, self.logger) - self.valves = DDValves(self.can_interface, self.logger) - self.voltages = DDVoltages(self.can_interface, self.logger) + self.alarms = DDAlarms(self.can_interface, self.logger) #: The Alarms module + self.balancing_chamber = DDBalancingChamber(self.can_interface, self.logger) #: The Balancing Chamber module + self.blood_leak = DDBloodLeak(self.can_interface, self.logger) #: The Blood Leak module + self.concentrate_pumps = DDConcentratePumps(self.can_interface, self.logger) #: The Concentrat Pumps module + self.conductivity_sensors = DDConductivitySensors(self.can_interface, self.logger) #: The Conductivity Sensors module + self.dialysate_pumps = DDDialysatePumps(self.can_interface, self.logger) #: The Dialysate Pumps module + self.drybicart = DDDryBicart(self.can_interface, self.logger) #: The Dry Bicarb module + self.dialysate_mixing = DDDialysateMixing(self.can_interface, self.logger) #: The Dialysate mixing module + self.events = DDEvents(self.can_interface, self.logger) #: The Events module + self.gen_dialysate = DDGenDialysate(self.can_interface, self.logger) #: The Generate Dialysate module + self.heaters = DDHeaters(self.can_interface, self.logger) #: The Heaters module + self.levels = DDLevels(self.can_interface, self.logger) #: The Levels module + self.post_gen_dialysate = DDPostGenDialysate(self.can_interface, self.logger) #: The Post Generate Dialysate module + self.pressure_sensors = DDPressureSensors(self.can_interface, self.logger) #: The Pressure Sensors module + self.pre_gen_dialysate = DDPreGenDialysate(self.can_interface, self.logger) #: The Pre Generate Dialysate module + self.rinse_pump = DDRinsePump(self.can_interface, self.logger) #: The Rinse Pump module + self.spent_chamber_fill = DDSpentChamberFill(self.can_interface, self.logger) #: The Spent Chamber module + self.substitution_pump = DDSubstitutionPump(self.can_interface, self.logger) #: The Substitution Pump module + self.temperature_sensors = DDTemperatureSensors(self.can_interface, self.logger) #: The Temperature Sensors module + self.test_configs = DDTestConfig(self.can_interface, self.logger) #: The Test Configs module + self.ultrafiltration = DDUltrafiltration(self.can_interface, self.logger) #: The Ultrafiltration module + self.valves = DDValves(self.can_interface, self.logger) #: The Valves module + self.voltages = DDVoltages(self.can_interface, self.logger) #: The Voltages module - self.ro_proxy = ROProxy(self.can_interface, self.logger) - self.td_proxy = TDProxy(self.can_interface, self.logger) + self.ro_proxy = ROProxy(self.can_interface, self.logger) #: The RO Proxy module (imitates commands sent by DD) + self.td_proxy = TDProxy(self.can_interface, self.logger) #: The TD Proxy module (imitates commands sent by UI) + def dd_set_logged_in_status(self, logged_in: bool = False): + """ + Callback for dd logged in status change. + + @param logged_in: Logged in status for DD + @return: None + """ + self.dd_logged_in = logged_in + + @publish(["msg_id_dd_debug_event", "dd_debug_events_timestamp","dd_debug_events"]) def _handler_dd_debug_event_sync(self, message, timestamp = 0.0): - payload = message['message'] message_length = payload[self._DD_DEBUG_EVENT_MSG_LEN_INDEX] temp_message = '' @@ -163,14 +179,6 @@ if self.dd_debug_event_index == self._DD_DEBUG_EVENT_LIST_COUNT: self.dd_debug_event_index = 0 - @publish(["dd_logged_in"]) - def dd_set_logged_in_status(self, logged_in: bool = False): - """ - Callback for dd logged in status change. - @param logged_in boolean logged in status for DD - @return: none - """ - self.dd_logged_in = logged_in @publish(["msg_id_dd_op_mode_data", "dd_op_mode_timestamp","dd_operation_mode", "dd_operation_sub_mode"]) def _handler_dd_op_mode_sync(self, message, timestamp = 0.0): @@ -181,15 +189,15 @@ @param message: published DD operation mode broadcast message @return: None """ - mode = struct.unpack('i', bytearray( - message['message'][MsgFieldPositions.START_POS_FIELD_1:MsgFieldPositions.END_POS_FIELD_1])) - smode = struct.unpack('i', bytearray( - message['message'][MsgFieldPositions.START_POS_FIELD_2:MsgFieldPositions.END_POS_FIELD_2])) + msg_list = [] + msg_list.append(('self.dd_operation_mode', DataTypes.U32)) + msg_list.append(('self.dd_operation_sub_mode', DataTypes.U32)) - self.dd_operation_mode = mode[0] - self.dd_operation_sub_mode = smode[0] + self.process_into_vars(decoder_list = msg_list, + message = message) self.dd_op_mode_timestamp = timestamp + @publish(["msg_id_dd_version_response", "dd_version, dd_fpga_version"]) def _handler_dd_version_response_sync(self,message, timestamp = 0.0): """ @@ -199,35 +207,52 @@ @return: None if not successful, the version string if unpacked successfully """ - major = struct.unpack(' 0 for each in [major, minor, micro, build, compatibility]]): - self.dd_version = f"v{major[0]}.{minor[0]}.{micro[0]}-{build[0]}.{compatibility[0]}" - self.logger.debug(f"DD VERSION: {self.dd_version}") + result = self.process_into_vars(decoder_list = msg_list, + message = message) + + if all([each is not None for each in [result['major'], result['minor'], result['micro'], result['build'], result['compatibility']]]): + self.dd_version = f"v{result['major']}.{result['minor']}.{result['micro']}-{result['build']}.{result['compatibility']}" + self.logger.debug(f'DD VERSION: {self.dd_version}') - if all([len(each) > 0 for each in [fpga_id, fpga_major, fpga_minor, fpga_lab]]): - self.dd_fpga_version = f"v{fpga_id[0]}.{fpga_major[0]}.{fpga_minor[0]}-{fpga_lab[0]}" - self.logger.debug(f"DD FPGA VERSION: {self.dd_fpga_version}") + if all([each is not None for each in [result['fpga_id'], result['fpga_major'], result['fpga_minor'], result['fpga_lab']]]): + self.dd_fpga_version = f"v{result['fpga_id']}.{result['fpga_major']}.{result['fpga_minor']}-{result['fpga_lab']}" + self.logger.debug(f'DD FPGA VERSION: {self.dd_fpga_version}') self.dd_version_response_timestamp = timestamp + + # def cmd_op_mode_broadcast_interval_override(self, ms: int, reset: int = NO_RESET) -> int: + # """ + # Constructs and sends the measured op mode broadcast interval override command + # Constraints: + # Must be logged into TD. + # Given interval must be non-zero and a multiple of the TD general task interval (50 ms). + + # @param ms: integer - interval (in ms) to override with + # @param reset: integer - 1 to reset a previous override, 0 to override + # @return: 1 if successful, zero otherwise + # """ + # return cmd_generic_broadcast_interval_override( + # ms = ms, + # reset = reset, + # channel_id = Channels.dialin_to_dd_ch_id, + # msg_id = MsgIds.MSG_ID_DD_OP_MODE_PUBLISH_INTERVAL_OVERRIDE_REQUEST, + # module_name = 'TD Operation Mode', + # logger = self.logger, + # can_interface = self.can_interface) + + def cmd_log_in_to_dd(self, resend: bool = False) -> int: """ Constructs and sends a login command via CAN bus. Login required before \n @@ -236,8 +261,7 @@ @param resend: (bool) if False (default), try to login once. Otherwise, tries to login indefinitely @return: 1 if logged in, 0 if log in failed """ - - message = DenaliMessage.build_message(channel_id=DenaliChannels.dialin_to_dd_ch_id, + message = CanMessage.build_message(channel_id=CanChannels.dialin_to_dd_ch_id, message_id=MsgIds.MSG_ID_DD_TESTER_LOGIN_REQUEST.value, payload=list(map(int, map(ord, self.DD_LOGIN_PASSWORD)))) @@ -247,18 +271,19 @@ received_message = self.can_interface.send(message, resend=resend) if received_message is not None: - if received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] == 1: + if received_message['message'][CanMessage.PAYLOAD_START_INDEX] == 1: self.logger.debug("Success: Logged In") self.dd_set_logged_in_status(True) #self._send_dd_checkin_message() # Timer starts interval first #self.can_interface.transmit_interval_dictionary[self.callback_id].start() else: self.logger.debug("Failure: Log In Failed.") - return received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] + return received_message['message'][CanMessage.PAYLOAD_START_INDEX] else: self.logger.debug("Login Timeout!!!!") return False + def cmd_dd_set_operation_mode(self, new_mode: int = 0) -> int: """ Constructs and sends a set operation mode request command via CAN bus. @@ -270,31 +295,20 @@ @param new_mode: ID of operation mode to transition to @return: 1 if successful, zero otherwise - """ - payload = integer_to_bytearray(new_mode) - message = DenaliMessage.build_message(channel_id=DenaliChannels.dialin_to_dd_ch_id, - message_id=MsgIds.MSG_ID_DD_SET_OPERATION_MODE_OVERRIDE_REQUEST.value, - payload=payload) + return cmd_generic_override( + payload = payload, + reset = NO_RESET, + channel_id = CanChannels.dialin_to_dd_ch_id, + msg_id = MsgIds.MSG_ID_DD_SET_OPERATION_MODE_OVERRIDE_REQUEST, + entity_name = 'DD Operation Mode', + override_text = dd_enum_repository.DDOpModes(new_mode).name, + logger = self.logger, + can_interface = self.can_interface) - self.logger.debug("Requesting DD mode change to " + str(new_mode)) - # Send message - received_message = self.can_interface.send(message) - - if received_message is not None: - if received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] == 1: - self.logger.debug("Success: Mode change accepted") - else: - self.logger.debug("Failure: Mode change rejected.") - return received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] - else: - self.logger.debug("DD mode change request Timeout!!!!") - return False - - def cmd_dd_software_reset_request(self) -> None: """ Constructs and sends an DD software reset request via CAN bus. @@ -303,18 +317,17 @@ @return: None """ + return cmd_generic_override( + payload = None, + reset = NO_RESET, + channel_id = CanChannels.dialin_to_dd_ch_id, + msg_id = MsgIds.MSG_ID_DD_SOFTWARE_RESET_REQUEST, + entity_name = 'DD Software Reset', + override_text = '', + logger = self.logger, + can_interface = self.can_interface) - message = DenaliMessage.build_message(channel_id=DenaliChannels.dialin_to_dd_ch_id, - message_id=MsgIds.MSG_ID_DD_SOFTWARE_RESET_REQUEST.value) - self.logger.debug("requesting DD software reset") - - # Send message - self.can_interface.send(message, 0) - self.logger.debug("Sent request to DD to reset...") - self.dd_set_logged_in_status(False) - - def cmd_dd_safety_shutdown_override(self, active: int, reset: int = NO_RESET) -> int: """ Constructs and sends an DD safety shutdown override command via CAN bus. @@ -324,28 +337,17 @@ @param active: int - True to activate safety shutdown, False to deactivate @param reset: integer - 1 to reset a previous override, 0 to override @return: 1 if successful, zero otherwise - """ - rst = integer_to_bytearray(reset) saf = integer_to_bytearray(active) payload = rst + saf - message = DenaliMessage.build_message(channel_id=DenaliChannels.dialin_to_dd_ch_id, - message_id=MsgIds.MSG_ID_DD_SAFETY_SHUTDOWN_OVERRIDE_REQUEST.value, - payload=payload) - - self.logger.debug("overriding FP safety shutdown") - - # Send message - received_message = self.can_interface.send(message) - - if received_message is not None: - if received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] == 1: - self.logger.debug("Safety shutdown signal overridden") - else: - self.logger.debug("Safety shutdown signal override failed.") - return received_message['message'][DenaliMessage.PAYLOAD_START_INDEX] - else: - self.logger.debug("Timeout!!!!") - return False + return cmd_generic_override( + payload = payload, + reset = NO_RESET, + channel_id = CanChannels.dialin_to_dd_ch_id, + msg_id = MsgIds.MSG_ID_DD_SET_OPERATION_MODE_OVERRIDE_REQUEST, + entity_name = 'DD Safety Shutdown', + override_text = str(active), + logger = self.logger, + can_interface = self.can_interface)