import struct import time from threading import Thread, Lock from typing import Optional import pandas as pd from nbus_api.nbus_sensor import NBusSensor from nbus_api.nbus_common_parser import NbusCommonParser from nbus_hal.crc8 import crc8 from nbus_hal.nbus_generic_port import * from nbus_types.nbus_address_type import NBusModuleAddress from nbus_types.nbus_data_fomat import NBusDataValue, NBusDataFormat from nbus_types.nbus_defines import * from nbus_types.nbus_exceptions.nbus_api_exception import NBusErrorAPI, NBusErrorAPIType from nbus_types.nbus_parameter_type import NBusParameterID, NBusParameterValue from nbus_types.nbus_status_type import NBusStatusType from nbus_types.nbus_sensor_count_type import NBusSensorCount from nbus_types.nbus_info_type import NBusModuleInfo from nbus_types.nbus_sensor_type import NBusSensorType @beartype class NBusSlaveModule: """ Class representing nBus slave module. """ def __init__(self, port: NBusPort, module_address: NBusModuleAddress, acquire_delay: float): """ Constructor. :param port: serial port :param module_address: address of module :param device_cnt: number of devices """ self.__port = port self.__module_addr = module_address self.__params = {} self.__devices = {} self._lock = Lock() # thread lock self.__in_acquisition = False # flag when in acquisition self.__acquire_delay = acquire_delay # intermediate delay between data fetching when data not ready self.__acquire_thread = None # thread for data acquisition self.__data_raw = bytearray() # raw data buffer self.__df = pd.DataFrame() # internal data frame self.__ts0 = None # 0-th timestamp self.__payload_size = 0 """ ================================================================================================================ Module General Methods ================================================================================================================ """ def init(self) -> None: """ Initialize the module from hardware. """ sensors = self.cmd_get_sensor_type() for sen_address, sen_type in sensors.items(): self.__devices[sen_address] = NBusSensor(self.__port, self.__module_addr, sen_address) self.__devices[sen_address].type = sen_type self.cmd_get_format() self.__calculate_payload_size() def get_devices(self) -> dict[NBusSensorAddress, NBusSensor]: """ Get module devices. :return: dictionary of connected devices """ return self.__devices """ ================================================================================================================ Module Get Commands ================================================================================================================ """ def cmd_get_echo(self, message: bytearray) -> bool: """ Get echo from module. :param message: message to send :return: status (True = echo, False = no echo) """ _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_ECHO, message) return response == list(message) def cmd_get_param(self, parameter: NBusParameterID) -> NBusParameterValue: """ Get single module parameter. :param parameter: parameter id :return: parameter value """ # get response resp_len, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_PARAM, bytearray([parameter.value])) # parse parameter param_id, param_val = NbusCommonParser.parameters_from_response(resp_len, response)[0] # store parameter value self.__params[param_id] = param_val return param_val def cmd_get_all_params(self) -> dict[NBusParameterID, NBusParameterValue]: """ Get all module parameters. :return: dict of parameter id and parameter value """ # get response resp_len, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_PARAM, bytearray([])) # parse parameters params = NbusCommonParser.parameters_from_response(resp_len, response) for param_id, param_val in params: # store parameters self.__params[param_id] = param_val return self.__params.copy() def cmd_get_sensor_cnt(self) -> NBusSensorCount: """ Get sensor count. :return: count of read-only and read-write sensors """ _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_SENSOR_CNT, bytearray([])) return NBusSensorCount(*response) def cmd_get_data(self) -> dict[NBusSensorAddress, list[NBusDataValue]]: """ Get data from all module sensors. :return: dict of device addresses and data values """ # get data resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_DATA, bytearray([])) # parse data begin_idx = 0 data = {} while begin_idx < resp_length: device_id = response[begin_idx] # handle errors if self.__devices[device_id].data_format is None: # check for format and params raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED) values, offset = NbusCommonParser.data_from_response(self.__devices[device_id].data_format, response[begin_idx:]) data[device_id] = values begin_idx += offset + 1 return data def cmd_get_info(self) -> NBusModuleInfo: """ Get module info. :return: module info """ response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_INFO, bytearray([])) name = str(response[1:9], "ascii") typ = str(response[9:12], "ascii") uuid = struct.unpack(" dict[NBusSensorAddress, NBusDataFormat]: """ Get format of all on-board sensors. :return: dict of sensor addresses and sensor formats """ # get response resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_FORMAT, bytearray([])) begin_idx = 0 formats = {} # parse format while begin_idx < resp_length: device_id = response[begin_idx] device_format = NbusCommonParser.format_from_response(response[begin_idx:begin_idx + 4]) self.__devices[device_id].data_format = device_format formats[device_id] = device_format begin_idx += 4 return formats def cmd_get_sensor_type(self) -> dict[NBusSensorAddress, NBusSensorType]: """ Get type of all on-board sensors. :return: dict of sensor addresses and sensor types """ # get response resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_SENSOR_TYPE, bytearray([])) # parse response types = {} i = 0 while i < resp_length - 1: types[NBusSensorAddress(response[i])] = NBusSensorType(response[i + 1]) i += 2 return types """ ================================================================================================================ Module Set Commands ================================================================================================================ """ def cmd_set_module_stop(self) -> None: """ Stop automatic measuring. :return: status """ self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_STOP, bytearray([])) def cmd_set_module_start(self) -> None: """ Start automatic measuring. :return: status """ self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_START, bytearray([])) def cmd_set_param(self, param: NBusParameterID, value: NBusParameterValue) -> NBusStatusType: """ Set module parameter. :param param: parameter ID :param value: parameter value :return: status """ # create request packet param_id_raw = struct.pack("B", param.value) param_val_raw = struct.pack(" dict[NBusParameterID, NBusStatusType]: """ Set multiple module parameters. :param params: parameters :return: dict od statuses """ # create request packet param_bytes = NbusCommonParser.parameters_to_request(params) # send request resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_PARAM, bytearray(param_bytes), long_answer=1.0) # parse statuses statuses = {} for i in range(0, resp_length - 1, 2): p_id = NBusParameterID(response[i]) if response[i+1] == NBusStatusType.STATUS_SUCCESS: self.__params[p_id] = params[p_id] # if success, store param statuses[p_id] = NBusStatusType(response[i + 1]) return statuses def cmd_set_calibrate(self) -> NBusStatusType: """ Send calibration command. :return: calibration status """ resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_CALIBRATE, bytearray([])) print(response) return NBusStatusType(response[0]) def cmd_set_data(self, data: dict[NBusSensorAddress, list[NBusDataValue]]) \ -> dict[NBusSensorAddress, NBusStatusType]: """ Set data to read-write sensors. :param data: :return: operation statuses """ # create request packet request = [] # transform data for addr in data.keys(): # handle errors if self.__devices[addr].data_format is None: # check for format and params raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED) raw_data = data[addr] request.append(addr) request.extend(NbusCommonParser.data_to_request(self.__devices[addr].data_format, raw_data)) # send request resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_DATA, bytearray(request)) # return response statuses statuses = {} for i in range(0, resp_length - 1, 2): statuses[NBusSensorAddress(response[i])] = NBusStatusType(response[i + 1]) return statuses """ ================================================================================================================ Automatic Data Stream Methods ================================================================================================================ """ def start_streaming(self) -> None: """ Start data streaming (e.g. auto-cast). """ self.cmd_set_module_start() self.__acquire_thread = Thread(target=self.__acquire_callback) self.__in_acquisition = True # end thread if running if self.__acquire_thread is not None and self.__acquire_thread.is_alive(): self.__acquire_thread.join() self.__acquire_thread.start() def stop_streaming(self): """ Stop data streaming (e.g. auto-cast). """ self.cmd_set_module_stop() self.__in_acquisition = False if self.__acquire_thread is not None and self.__acquire_thread.is_alive(): self.__acquire_thread.join() self.__port.flush() def fetch_stream_chunk(self) -> pd.DataFrame: """ Fetch data from stream (e.g. auto-cast). Can be called anytime. It not erase internal dataframe. :return: stream data frame """ with self._lock: packets = self.__data_raw.split(NBUS_BRIDGE_DATA_HDR)[1:] packet_cnt = len(packets) - self.__in_acquisition parsed_packets = [] # parse packets for i in range(packet_cnt): data = self._parse_data_from_stream_packet(packets[i]) if data is not None: parsed_packets.append(data) else: print("damaged: ", i, packets[i]) # extend internal dataframe if parsed_packets: data_frame = pd.DataFrame(parsed_packets) self.__df = pd.concat([self.__df, data_frame], ignore_index=True).copy(deep=True) else: data_frame = pd.DataFrame() # erase raw data buffer if self.__in_acquisition: unparsed_bytes = len(packets[-1]) + NBUS_BRIDGE_DATA_HDR_SIZE self.__data_raw = self.__data_raw[-unparsed_bytes:] else: self.__data_raw = bytearray() self._transform_timestamp(data_frame) return data_frame def fetch_full_stream(self) -> pd.DataFrame: """ Fetch all data from stream (e.g. auto-cast). Must be called after stop_streaming() method. It will erase internal dataframe. :return: stream dataframe """ if self.__in_acquisition: return pd.DataFrame() self.fetch_stream_chunk() df = self.__df self._transform_timestamp(df) self.__df = pd.DataFrame() self.__ts0 = None return df """ ================================================================================================================ Internal Helper Methods ================================================================================================================ """ def _transform_timestamp(self, data_frame: pd.DataFrame) -> None: """ Transform timestamp values in dataframe. :param data_frame: dataframe to transform """ if not data_frame.empty and "TS" in self.__df.columns: if self.__ts0 is None: self.__ts0 = self.__df["TS"].iloc[0] data_frame["TS"] -= self.__ts0 def _parse_data_from_stream_response(self, module_address: NBusModuleAddress, resp_length: int, response: bytearray) \ -> dict[str, NBusDataValue]: """ Parse data of slave from stream response. :param module_address: address of module :param resp_length: length of response :param response: raw data :return: dict of data values """ data_offset = 0 data = {} while data_offset < resp_length: device_id = response[data_offset] if self.__devices[device_id].data_format is None: raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED) values, offset = NbusCommonParser.data_from_response(self.__devices[device_id].data_format, response[data_offset:]) data_tag = str(module_address) + "." + str(device_id) for i in range(len(values)): data[data_tag + "." + str(i + 1)] = values[i] data_offset += offset + NBUS_SA_SIZE return data def _parse_data_from_stream_packet(self, data_packet: bytearray) \ -> Optional[dict[str, NBusDataValue]]: """ Parse the stream-mode data packet. :param data_packet: packet to parse :return: dictionary of data values or None """ packet_size = len(data_packet) packet_crc = crc8(data_packet[:NBUS_CRC_ADDR]) # check validity if packet_size < NBUS_TS_SIZE + NBUS_CRC_SIZE or data_packet[NBUS_CRC_ADDR] != packet_crc: return None # parse data try: data_offset = 0 ts = struct.unpack(" None: """ Thread worker for data acquisition. """ while self.__in_acquisition: data = self.__port.try_read() if not data: time.sleep(self.__acquire_delay) continue with self._lock: self.__data_raw.extend(data)