| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543 |
- 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("<I", bytearray(response[12:16]))[0]
- hw = str(response[16:19], "ascii")
- fw = str(response[19:22], "ascii")
- mem_id = struct.unpack("<Q", bytearray(response[22:30]))[0]
- ro_count = int(response[30])
- rw_count = int(response[31])
- return NBusModuleInfo(module_name=name, module_type=typ, uuid=uuid, hw=hw, fw=fw, memory_id=mem_id,
- read_only_sensors=ro_count, read_write_sensors=rw_count)
- def cmd_get_format(self) -> 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("<I",value)
- param_bytes = bytearray(param_id_raw) + bytearray(param_val_raw)
- # proceed request
- _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_PARAM, param_bytes)
- # if response is valid, store parameter
- if response[0] == param.value and response[1] == NBusStatusType.STATUS_SUCCESS:
- self.__params[param] = value
- return NBusStatusType(response[1])
- def cmd_set_multi_params(self, params: dict[NBusParameterID, NBusParameterValue]) \
- -> 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("<I", data_packet[data_offset: data_offset + NBUS_TS_SIZE])[0]
- data = {"TS": ts}
- data_offset += NBUS_TS_SIZE
- module_addr = data_packet[data_offset]
- data_offset += NBUS_MA_SIZE
- packet = data_packet[data_offset: data_offset + self.__payload_size]
- data |= self._parse_data_from_stream_response(module_addr, self.__payload_size, packet)
- data_offset += self.__payload_size
- return data
- except Exception:
- return None
- def __calculate_payload_size(self):
- """
- Calculate payload size of stream-mode packet.
- :return: payload size in bytes
- """
- for device in self.__devices.values():
- fmt = device.data_format
- self.__payload_size += fmt.byte_length * fmt.samples + NBUS_SA_SIZE
- def __acquire_callback(self) -> 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)
|