nbus_bridge.py 7.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230
  1. import struct
  2. from collections import deque
  3. from dataclasses import dataclass
  4. from nbus_api.nbus_common_parser import NbusCommonParser
  5. from nbus_api.nbus_module_slave import NBusSlaveModule
  6. from nbus_api.nbus_sensor import NBusSensor
  7. from nbus_hal.nbus_serial.serial_port import *
  8. from nbus_types.nbus_exceptions.nbus_api_exception import NBusErrorAPI, NBusErrorAPIType
  9. from nbus_types.nbus_sensor_count_type import NBusSensorCount
  10. import time
  11. from threading import Thread, Lock
  12. from collections import namedtuple
  13. import pandas as pd
  14. NBUS_RX_META = 4
  15. NBUS_FMT_SIZE = 4
  16. NBUS_TS_SIZE = 4
  17. NBUS_CRC_SIZE = 1
  18. NBUS_MA_SIZE = 1
  19. NBUS_SA_SIZE = 1
  20. NBUS_CRC_ADDR = -1
  21. NBUS_BRIDGE_DATA_HDR = bytearray([0x00] + [0xFF] * 8 + [0x00])
  22. @dataclass
  23. class NBusSlaveMeta:
  24. obj: NBusSlaveModule
  25. cnt: int
  26. packet_size: int
  27. @beartype
  28. class NBusBridge:
  29. def __init__(self, serial_port: NBusSerialPort):
  30. """
  31. Constructor.
  32. :param serial_port: serial port
  33. """
  34. self.__port = serial_port
  35. self.__slaves_meta = {}
  36. self.__slaves = {}
  37. self.__slaves_meta = {}
  38. self.__scan_thread = None
  39. self.__data_raw = bytearray([])
  40. self.buf = deque()
  41. self.lock = Lock()
  42. self.__in_scan = False
  43. self.__packet_size = 0
  44. self.df = pd.DataFrame()
  45. def init_from_network(self):
  46. self.cmd_get_slaves()
  47. self.cmd_get_format()
  48. def cmd_get_slaves(self):
  49. resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_SLAVES, bytearray([]))
  50. slaves = {}
  51. data_offset = 0
  52. while data_offset < resp_length - 2:
  53. slave_addr = response[data_offset]
  54. slave_sensor_cnt = NBusSensorCount(response[data_offset + 1] , response[data_offset + 2])
  55. slaves[slave_addr] = slave_sensor_cnt
  56. self.__slaves_meta[slave_addr] = NBusSlaveMeta(NBusSlaveModule(self.__port, slave_addr), slave_sensor_cnt, 0)
  57. data_offset += 3
  58. return slaves
  59. def cmd_get_format(self):
  60. resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_FORMAT, bytearray([]))
  61. data_offset = 0
  62. fmt = {}
  63. while data_offset < resp_length:
  64. slave_addr = response[data_offset]
  65. data_offset += NBUS_MA_SIZE
  66. slave_sensor_cnt = self.__slaves_meta[slave_addr].cnt
  67. format_len = NBUS_FMT_SIZE * (slave_sensor_cnt.read_only_count + slave_sensor_cnt.read_write_count)
  68. fmt |= self._set_slave_format_from_response(slave_addr, format_len,
  69. response[data_offset: data_offset + format_len])
  70. data_offset += format_len
  71. self._set_data_packet_size()
  72. return fmt
  73. def cmd_get_data(self):
  74. resp_length, *response = self.__port.request_bridge(NBusCommand.CMD_GET_FORMAT, bytearray([]))
  75. def scan_fn(self):
  76. while self.__in_scan:
  77. with self.lock:
  78. if self.__port.get_port().in_waiting > 0:
  79. self.__data_raw.extend(self.__port.get_port().read())
  80. else:
  81. time.sleep(0.05)
  82. def start_data_streaming(self):
  83. self.__port.send_bridge(NBusCommand.CMD_SET_START, bytearray([]))
  84. self.__scan_thread = Thread(target=self.scan_fn)
  85. self.__in_scan = True
  86. self.__scan_thread.start()
  87. def stop_data_streaming(self):
  88. self.__port.send_bridge(NBusCommand.CMD_SET_STOP, bytearray([]))
  89. self.__in_scan = False
  90. self.__scan_thread.join()
  91. while self.__port.get_port().in_waiting:
  92. time.sleep(0.05)
  93. self.__port.get_port().reset_input_buffer()
  94. print("flushed baby")
  95. def get_data(self):
  96. return self.__data_raw
  97. def dataen(self):
  98. with (self.lock):
  99. frames = self.__data_raw.split(NBUS_BRIDGE_DATA_HDR)
  100. contain_end = int(self.__in_scan)
  101. frames_cnt = len(frames) - contain_end
  102. for i in range(frames_cnt):
  103. frame = frames[i]
  104. frame_len = len(frame)
  105. frame_crc = crc8(NBUS_BRIDGE_DATA_HDR + frame[:NBUS_CRC_ADDR])
  106. data_offset = 0
  107. if frame_len == self.__packet_size and frame[NBUS_CRC_ADDR] == frame_crc:
  108. ts = struct.unpack("<I", frame[data_offset : data_offset + NBUS_TS_SIZE])[0]
  109. data = {"TS": ts}
  110. data_offset += NBUS_TS_SIZE
  111. while data_offset < frame_len - 1:
  112. module_addr = frame[data_offset]
  113. data_offset += NBUS_MA_SIZE
  114. packet_size = self.__slaves_meta[module_addr].packet_size
  115. packet = frame[data_offset : data_offset + packet_size]
  116. data |= self._get_data_from_response(module_addr, packet_size, packet)
  117. data_offset += packet_size
  118. self.df = pd.concat([self.df, pd.DataFrame([data])])
  119. else:
  120. print("Zahodeny")
  121. if contain_end:
  122. unparsed_bytes = len(frames[-1]) + len(NBUS_BRIDGE_DATA_HDR)
  123. self.__data_raw = self.__data_raw[-unparsed_bytes:]
  124. return []
  125. def _set_data_packet_size(self):
  126. self.__packet_size = NBUS_TS_SIZE + NBUS_CRC_SIZE
  127. for slave_meta in self.__slaves_meta.values():
  128. packet_size = 0
  129. for device in slave_meta.obj.get_devices().values():
  130. fmt = device.data_format
  131. packet_size += fmt.byte_length * fmt.samples + NBUS_SA_SIZE
  132. slave_meta.packet_size = packet_size
  133. self.__packet_size += packet_size + NBUS_MA_SIZE
  134. def _set_slave_format_from_response(self, slave_addr, resp_length, response):
  135. data_offset = 0
  136. formats = {}
  137. devices = self.__slaves_meta[slave_addr].obj.get_devices()
  138. # parse format
  139. while data_offset < resp_length:
  140. device_id = response[data_offset]
  141. device_format = NbusCommonParser.format_from_response(response[data_offset:data_offset + 4])
  142. if device_id not in devices:
  143. devices[device_id] = NBusSensor(device_id)
  144. devices[device_id].data_format = device_format
  145. tag = str(slave_addr) + "." + str(device_id)
  146. formats[tag] = device_format
  147. data_offset += NBUS_FMT_SIZE
  148. return formats
  149. def _get_data_from_response(self, slave_addr, resp_length, response):
  150. """
  151. Get data from all module sensors.
  152. :return: dict of device addresses and data values
  153. """
  154. # parse data
  155. data_offset = 0
  156. data = {}
  157. devices = self.__slaves_meta[slave_addr].obj.get_devices()
  158. while data_offset < resp_length:
  159. device_id = response[data_offset]
  160. # handle errors
  161. if devices[device_id].data_format is None: # check for format and params
  162. raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED)
  163. values, offset = NbusCommonParser.data_from_response(devices[device_id].data_format,
  164. response[data_offset:])
  165. tag = str(slave_addr) + "." + str(device_id)
  166. for i in range(len(values)):
  167. data[tag + "." + str(i + 1)] = values[i]
  168. data_offset += offset + NBUS_SA_SIZE
  169. return data