import struct from collections import deque from dataclasses import dataclass from nbus_api.nbus_common_parser import NbusCommonParser from nbus_api.nbus_module_slave import NBusSlaveModule from nbus_api.nbus_sensor import NBusSensor from nbus_hal.nbus_serial.serial_port import * from nbus_types.nbus_exceptions.nbus_api_exception import NBusErrorAPI, NBusErrorAPIType from nbus_types.nbus_sensor_count_type import NBusSensorCount import time from threading import Thread, Lock from collections import namedtuple import pandas as pd NBUS_RX_META = 4 NBUS_FMT_SIZE = 4 NBUS_TS_SIZE = 4 NBUS_CRC_SIZE = 1 NBUS_MA_SIZE = 1 NBUS_SA_SIZE = 1 NBUS_CRC_ADDR = -1 NBUS_BRIDGE_DATA_HDR = bytearray([0x00] + [0xFF] * 8 + [0x00]) @dataclass class NBusSlaveMeta: obj: NBusSlaveModule cnt: int packet_size: int @beartype class NBusBridge: def __init__(self, serial_port: NBusSerialPort): """ Constructor. :param serial_port: serial port """ self.__port = serial_port self.__slaves_meta = {} self.__slaves = {} self.__slaves_meta = {} self.__scan_thread = None self.__data_raw = bytearray([]) self.buf = deque() self.lock = Lock() self.__in_scan = False self.__packet_size = 0 self.df = pd.DataFrame() def init_from_network(self): self.cmd_get_slaves() self.cmd_get_format() def cmd_get_slaves(self): resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_SLAVES, bytearray([])) slaves = {} data_offset = 0 while data_offset < resp_length - 2: slave_addr = response[data_offset] slave_sensor_cnt = NBusSensorCount(response[data_offset + 1] , response[data_offset + 2]) slaves[slave_addr] = slave_sensor_cnt self.__slaves_meta[slave_addr] = NBusSlaveMeta(NBusSlaveModule(self.__port, slave_addr), slave_sensor_cnt, 0) data_offset += 3 return slaves def cmd_get_format(self): resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_FORMAT, bytearray([])) data_offset = 0 fmt = {} while data_offset < resp_length: slave_addr = response[data_offset] data_offset += NBUS_MA_SIZE slave_sensor_cnt = self.__slaves_meta[slave_addr].cnt format_len = NBUS_FMT_SIZE * (slave_sensor_cnt.read_only_count + slave_sensor_cnt.read_write_count) fmt |= self._set_slave_format_from_response(slave_addr, format_len, response[data_offset: data_offset + format_len]) data_offset += format_len self._set_data_packet_size() return fmt def cmd_get_data(self): resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_FORMAT, bytearray([])) def scan_fn(self): while self.__in_scan: with self.lock: if self.__port.get_port().in_waiting > 0: self.__data_raw.extend(self.__port.get_port().read()) else: time.sleep(0.05) def start_data_streaming(self): self.__port.send_bridge(NBusCommand.CMD_SET_START, bytearray([])) self.__scan_thread = Thread(target=self.scan_fn) self.__in_scan = True self.__scan_thread.start() def stop_data_streaming(self): self.__port.send_bridge(NBusCommand.CMD_SET_STOP, bytearray([])) self.__in_scan = False self.__scan_thread.join() while self.__port.get_port().in_waiting: time.sleep(0.05) self.__port.get_port().reset_input_buffer() print("flushed baby") def get_data(self): return self.__data_raw def dataen(self): with (self.lock): frames = self.__data_raw.split(NBUS_BRIDGE_DATA_HDR) contain_end = int(self.__in_scan) frames_cnt = len(frames) - contain_end for i in range(frames_cnt): frame = frames[i] frame_len = len(frame) frame_crc = crc8(NBUS_BRIDGE_DATA_HDR + frame[:NBUS_CRC_ADDR]) data_offset = 0 if frame_len == self.__packet_size and frame[NBUS_CRC_ADDR] == frame_crc: ts = struct.unpack("