| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230 |
- 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("<I", frame[data_offset : data_offset + NBUS_TS_SIZE])[0]
- data = {"TS": ts}
- data_offset += NBUS_TS_SIZE
- while data_offset < frame_len - 1:
- module_addr = frame[data_offset]
- data_offset += NBUS_MA_SIZE
- packet_size = self.__slaves_meta[module_addr].packet_size
- packet = frame[data_offset : data_offset + packet_size]
- data |= self._get_data_from_response(module_addr, packet_size, packet)
- data_offset += packet_size
- self.df = pd.concat([self.df, pd.DataFrame([data])])
- else:
- print("Zahodeny")
- if contain_end:
- unparsed_bytes = len(frames[-1]) + len(NBUS_BRIDGE_DATA_HDR)
- self.__data_raw = self.__data_raw[-unparsed_bytes:]
- return []
- def _set_data_packet_size(self):
- self.__packet_size = NBUS_TS_SIZE + NBUS_CRC_SIZE
- for slave_meta in self.__slaves_meta.values():
- packet_size = 0
- for device in slave_meta.obj.get_devices().values():
- fmt = device.data_format
- packet_size += fmt.byte_length * fmt.samples + NBUS_SA_SIZE
- slave_meta.packet_size = packet_size
- self.__packet_size += packet_size + NBUS_MA_SIZE
- def _set_slave_format_from_response(self, slave_addr, resp_length, response):
- data_offset = 0
- formats = {}
- devices = self.__slaves_meta[slave_addr].obj.get_devices()
- # parse format
- while data_offset < resp_length:
- device_id = response[data_offset]
- device_format = NbusCommonParser.format_from_response(response[data_offset:data_offset + 4])
- if device_id not in devices:
- devices[device_id] = NBusSensor(device_id)
- devices[device_id].data_format = device_format
- tag = str(slave_addr) + "." + str(device_id)
- formats[tag] = device_format
- data_offset += NBUS_FMT_SIZE
- return formats
- def _get_data_from_response(self, slave_addr, resp_length, response):
- """
- Get data from all module sensors.
- :return: dict of device addresses and data values
- """
- # parse data
- data_offset = 0
- data = {}
- devices = self.__slaves_meta[slave_addr].obj.get_devices()
- while data_offset < resp_length:
- device_id = response[data_offset]
- # handle errors
- if devices[device_id].data_format is None: # check for format and params
- raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED)
- values, offset = NbusCommonParser.data_from_response(devices[device_id].data_format,
- response[data_offset:])
- tag = str(slave_addr) + "." + str(device_id)
- for i in range(len(values)):
- data[tag + "." + str(i + 1)] = values[i]
- data_offset += offset + NBUS_SA_SIZE
- return data
|