########################################################################### # # Copyright (c) 2020-2025 Diality Inc. - All Rights Reserved. # # THIS CODE MAY NOT BE COPIED OR REPRODUCED IN ANY FORM, IN PART OR IN # WHOLE, WITHOUT THE EXPLICIT PERMISSION OF THE COPYRIGHT OWNER. # # @file treatment.py # # @author (last) Zoltan Miskolci # @date (last) 05-May-2026 # @author (original) Michael Garthwaite # @date (original) 22-Apr-2025 # ############################################################################ # Module imports from logging import Logger # Project imports from leahi_dialin.common.constants import NO_RESET from leahi_dialin.common.generic_defs import DataTypes from leahi_dialin.common.msg_ids import MsgIds from leahi_dialin.common.override_templates import cmd_generic_override from leahi_dialin.common import td_enum_repository from leahi_dialin.protocols.CAN import CanMessenger, CanChannels from leahi_dialin.utils.abstract_classes import AbstractSubSystem from leahi_dialin.utils.base import publish from leahi_dialin.utils.conversions import integer_to_bytearray, float_to_bytearray class TDTreatment(AbstractSubSystem): """ Treatment Delivery (TD) Dialin API sub-class for treatment related commands. """ def __init__(self, can_interface: CanMessenger, logger: Logger): """ TDTreatment constructor """ super().__init__() self.can_interface = can_interface self.logger = logger if self.can_interface is not None: self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_TREATMENT_PARAM_RANGES.value, function = self._handler_treatment_param_ranges_sync) self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_SALINE_BOLUS_DATA.value, function = self._handler_saline_bolus_sync) self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_ULTRAFILTRATION_DATA.value, function = self._handler_uf_sync) self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_TREATMENT_TIME_DATA.value, function = self._handler_treatment_time_sync) self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_TREATMENT_STATE_DATA.value, function = self._handler_treatment_state_sync) self.can_interface.register_receiving_publication_function(channel_id = CanChannels.td_sync_broadcast_ch_id, message_id = MsgIds.MSG_ID_TD_RSP_CURRENT_TREATMENT_PARAMETERS.value, function = self._handler_resp_treatment_parameters_sync) # Treatment param ranges data self.tx_params_timestamp = 0 #: The timestamp of the latest Treatment Parameters message self.treatment_param_ranges = { 'min_tx_time': 0, # The Minimum Treatment time 'max_tx_time': 0, # The Maximum Treatment time 'min_uf_volume': 0.0, # The Minimum Ultrafiltration volume 'max_uf_volume': 0.0, # The Maximum Ultrafiltration volume 'min_dial_rate': 0, # The Minimum Dialysate rate 'max_dial_rate': 0 # The Maximum Dialysate rate } # Saline Bolus data self.saline_bolus_timestamp = 0 #: The timestamp of the latest Saline Bolus message self.saline_bolus = { 'tgt_saline_volume': 0, # Target Saline Bolus volume in mL 'cum_saline_volume': 0.0, # Cum Saline Bolus volume in mL 'bol_saline_volume': 0.0, # Bol Saline Bolus volume in mL 'saline_bolus_state': 0 # The State of the Saline Bolus } # Saline Bolus data self.uf_timestamp = 0 #: The timestamp of the latest Ultrafiltration message self.ultrafiltration = { 'set_uf_volume': 0.0, # The Ultrafiltration volume in L 'tgt_uf_rate': 0.0, # The Ultrafiltration rate in L/hr 'uf_volume_delivered': 0.0, # How much was delivered during Ultrafiltration in L 'uf_state': 0 # The State of the Ultrafiltration } # Treatment Time Data self.tx_time_timestamp = 0 #: The timestamp of the latest Treatment message self.treatment_times ={ 'tx_time_prescribed': 0, # Total Treatment duration 'tx_time_elapsed': 0, # How much time is elapsed in the Treatment 'tx_time_remaining': 0 # How much time is remaining from the Treatment } # Treatment State Data self.tx_state_timestamp = 0 #: The timestamp of the latest Treatment State message self.treatment_states = { 'tx_sub_mode': 0, # The Treatment Operation Sub-Mode 'blood_prime_state': 0, # The Blood Prime Operation Sub-Mode 'dialysis_state': 0, # The Dialysis Operation Sub-Mode 'isolated_uf_state': 0, # The Isolated Ultrafiltration Operation Sub-Mode 'tx_stop_state': 0, # The Treatment Stop Operation Sub-Mode 'rinseback_state': 0, # The Rinsback Operation Sub-Mode 'tx_recirc_state': 0, # The Recirculation Operation Sub-Mode 'tx_end_state': 0 # The Treatment End Operation Sub-Mode } # Treatment Parameters Data. Most recent response. self.tx_param_req_timestamp = 0 #: The timestamp of the latest Treatment Parameters Request message self.treatment_paramters = { 'blood_flow_rate': 0, # The Blood flow rate 'dialysate_flow_rate': 0, # The Dialysate flow rate 'tx_duration': 0, # The Treatment duration 'saline_bolus_volume': 0, # The Saline Bolus volume 'hep_stop_time': 0, # The Heparin stop time 'hep_time': 0, # The Heparin start time 'acid_con': 0, # The selected Acid concentrate's option index 'bicarb_con': 0, # The selected Bicarb concentrate's option index 'dialyzer_type': 0, # The selected Dialyser option index 'bp_interval': 0, # The Body Pulse interval 'rb_flow_rate': 0, # The RB flow rate 'rb_volume': 0, # The RB volume 'art_pressure_window': 0, # The Artery pressure window duration 'venous_pressure_window': 0, # The Venous pressure window duration 'venous_asymm_window': 0, # The Venous asymmetric window duration 'tmp_limit_window': 0, # The Transmembrane limit window duration 'dialysate_temp': 0, # The Dialysate temperature 'hep_dispense_rate': 0, # The Heparin dispense rate 'hep_bolus_vol': 0, # The Heparin bolus volume 'uf_vol': 0 # The Ultrafiltration volume } # ============================================================ Properties ============================================================ @property def treatment_param_ranges(self) -> dict: """ The Treatment Parameter limits """ return self._min_tx_time @treatment_param_ranges.setter def treatment_param_ranges(self, value): self._min_tx_time = value @property def saline_bolus(self) -> dict: """ The Saline Bolus data """ return self._saline_bolus @saline_bolus.setter def saline_bolus(self, value): self._saline_bolus = value @property def ultrafiltration(self) -> int: """ The Ultrafiltration data """ return self._ultrafiltration @ultrafiltration.setter def ultrafiltration(self, value): self._ultrafiltration = value @property def treatment_times(self) -> int: """ The Treatment times """ return self._treatment_times @treatment_times.setter def treatment_times(self, value): self._treatment_times = value @property def treatment_states(self) -> int: """ The Treatment states """ return self._treatment_states @treatment_states.setter def treatment_states(self, value): self._treatment_states = value @property def treatment_paramters(self) -> int: """ The Treatment Parameters """ return self._treatment_paramters @treatment_paramters.setter def treatment_paramters(self, value): self._treatment_paramters = value # ============================================================ Handlers ============================================================ @publish(["msg_id_td_treatment_param_ranges", "min_tx_time","max_tx_time","min_uf_volume","max_uf_volume", "min_dial_rate","max_dial_rate","tx_params_timestamp"]) def _handler_treatment_param_ranges_sync(self, message, timestamp=0.0): """ Handles published treatment parameter range data messages. @param message: published treatment parameter range data message @return: none """ msg_list = [] msg_list.append((self.treatment_param_ranges, 'min_tx_time', DataTypes.U32)) msg_list.append((self.treatment_param_ranges, 'max_tx_time', DataTypes.U32)) msg_list.append((self.treatment_param_ranges, 'min_uf_volume', DataTypes.F32)) msg_list.append((self.treatment_param_ranges, 'max_uf_volume', DataTypes.F32)) msg_list.append((self.treatment_param_ranges, 'min_dial_rate', DataTypes.U32)) msg_list.append((self.treatment_param_ranges, 'max_dial_rate', DataTypes.U32)) self.process_into_vars(decoder_list = msg_list, message = message) self.tx_params_timestamp = timestamp @publish(["msg_id_td_saline_bolus_data", "tgt_saline_volume","cum_saline_volume","bol_saline_volume", "saline_bolus_state","saline_bolus_timestamp"]) def _handler_saline_bolus_sync(self, message, timestamp=0.0): """ Handles published saline bolus data messages. @param message: published saline bolus data message @return: none """ msg_list = [] msg_list.append((self.saline_bolus, 'tgt_saline_volume', DataTypes.U32)) msg_list.append((self.saline_bolus, 'cum_saline_volume', DataTypes.F32)) msg_list.append((self.saline_bolus, 'bol_saline_volume', DataTypes.F32)) msg_list.append((self.saline_bolus, 'saline_bolus_state', DataTypes.U32)) self.process_into_vars(decoder_list = msg_list, message = message) self.saline_bolus_timestamp = timestamp @publish(["msg_id_td_uf_data", "set_uf_volume","tgt_uf_rate","uf_volume_delivered", "uf_state","uf_timestamp"]) def _handler_uf_sync(self, message, timestamp=0.0): """ Handles published ultrafiltration data messages. @param message: published ultrafiltration data message @return: none """ msg_list = [] msg_list.append((self.ultrafiltration, 'set_uf_volume', DataTypes.F32)) msg_list.append((self.ultrafiltration, 'tgt_uf_rate', DataTypes.F32)) msg_list.append((self.ultrafiltration, 'uf_volume_delivered', DataTypes.F32)) msg_list.append((self.ultrafiltration, 'uf_state', DataTypes.U32)) self.process_into_vars(decoder_list = msg_list, message = message) self.uf_timestamp = timestamp @publish(["msg_id_td_treatment_time_data", "tx_time_prescribed","tx_time_elapsed","tx_time_remaining", "tx_time_timestamp"]) def _handler_treatment_time_sync(self, message, timestamp=0.0): """ Handles published treatment time data messages. @param message: published treatment time data message @return: none """ msg_list = [] msg_list.append((self.treatment_times, 'tx_time_prescribed', DataTypes.U32)) msg_list.append((self.treatment_times, 'tx_time_elapsed', DataTypes.U32)) msg_list.append((self.treatment_times, 'tx_time_remaining', DataTypes.U32)) self.process_into_vars(decoder_list = msg_list, message = message) self.tx_time_timestamp = timestamp @publish(["msg_id_td_treatment_state_data", "tx_sub_mode","blood_prime_state","dialysis_state","isolated_uf_state", "tx_stop_state","rinseback_state","tx_recirc_state","tx_end_state","tx_state_timestamp"]) def _handler_treatment_state_sync(self, message, timestamp=0.0): """ Handles published treatment state data messages. @param message: published treatment state data message @return: none """ msg_list = [] msg_list.append((self.treatment_states, 'tx_sub_mode', DataTypes.U32)) msg_list.append((self.treatment_states, 'blood_prime_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'dialysis_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'isolated_uf_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'tx_stop_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'rinseback_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'tx_recirc_state', DataTypes.U32)) msg_list.append((self.treatment_states, 'tx_end_state', DataTypes.U32)) self.process_into_vars(decoder_list = msg_list, message = message) self.tx_state_timestamp = timestamp @publish(["msg_id_td_rsp_current_treatment_parameters", "blood_flow_rate", "dialysate_flow_rate", "tx_duration", "saline_bolus_volume", "hep_stop_time", "hep_time", "acid_con", "bicarb_con", "dialyzer_type", "bp_interval", "rb_flow_rate", "rb_volume", "art_pressure_window","venous_pressure_window","venous_asymm_window","tmp_limit_window", "dialysate_temp", "hep_dispense_rate", "hep_bolus_vol", "uf_vol","tx_param_req_timestamp"]) def _handler_resp_treatment_parameters_sync(self, message, timestamp=0.0): """ Handles published treatment parameter response data messages. @param message: published treatment parameter response data message @return: none """ msg_list = [] msg_list.append((self.treatment_paramters, 'blood_flow_rate', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'dialysate_flow_rate', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'tx_duration', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'saline_bolus_volume', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'hep_stop_time', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'hep_time', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'acid_con', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'bicarb_con', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'dialyzer_type', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'bp_interval', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'rb_flow_rate', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'rb_volume', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'art_pressure_window', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'venous_pressure_window', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'venous_asymm_window', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'tmp_limit_window', DataTypes.U32)) msg_list.append((self.treatment_paramters, 'dialysate_temp', DataTypes.F32)) msg_list.append((self.treatment_paramters, 'hep_dispense_rate', DataTypes.F32)) msg_list.append((self.treatment_paramters, 'hep_bolus_vol', DataTypes.F32)) msg_list.append((self.treatment_paramters, 'uf_vol', DataTypes.F32)) self.process_into_vars(decoder_list = msg_list, message = message) self.tx_param_req_timestamp = timestamp # ============================================================ Overrides and Requests ============================================================ def cmd_set_treatment_parameter(self, tx_param_id: int = 0, tx_param_value = 0 ): """ Constructs and sends set treatment parameter command to the TD. Constraints: Must be logged into TD. @param tx_param_value: varied - value to set the treatment parameter @param tx_param_id: integer - the param id of the treatment paramter @return: 1 if successful, zero otherwise """ if tx_param_id <= td_enum_repository.TDTreatmentParameters.TREATMENT_PARAM_RINSEBACK_VOLUME.value: tpv = integer_to_bytearray(tx_param_value) elif tx_param_id >= td_enum_repository.TDTreatmentParameters.TREATMENT_PARAM_DIALYSATE_TEMPERATURE.value: tpv = float_to_bytearray(tx_param_value) else: tpv = integer_to_bytearray(tx_param_value) idx = integer_to_bytearray(tx_param_id) payload = idx + tpv param_name = td_enum_repository.TDTreatmentParameters(tx_param_id).name return cmd_generic_override( payload = payload, reset = NO_RESET, channel_id = CanChannels.dialin_to_td_ch_id, msg_id = MsgIds.MSG_ID_TD_SET_TREATMENT_PARAMETER, entity_name = f'TD {param_name}', override_text = str(tx_param_value), logger = self.logger, can_interface = self.can_interface)