nbus_module_slave.py 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543
  1. import struct
  2. import time
  3. from threading import Thread, Lock
  4. from typing import Optional
  5. import pandas as pd
  6. from nbus_api.nbus_sensor import NBusSensor
  7. from nbus_api.nbus_common_parser import NbusCommonParser
  8. from nbus_hal.crc8 import crc8
  9. from nbus_hal.nbus_generic_port import *
  10. from nbus_types.nbus_address_type import NBusModuleAddress
  11. from nbus_types.nbus_data_fomat import NBusDataValue, NBusDataFormat
  12. from nbus_types.nbus_defines import *
  13. from nbus_types.nbus_exceptions.nbus_api_exception import NBusErrorAPI, NBusErrorAPIType
  14. from nbus_types.nbus_parameter_type import NBusParameterID, NBusParameterValue
  15. from nbus_types.nbus_status_type import NBusStatusType
  16. from nbus_types.nbus_sensor_count_type import NBusSensorCount
  17. from nbus_types.nbus_info_type import NBusModuleInfo
  18. from nbus_types.nbus_sensor_type import NBusSensorType
  19. @beartype
  20. class NBusSlaveModule:
  21. """
  22. Class representing nBus slave module.
  23. """
  24. def __init__(self, port: NBusPort, module_address: NBusModuleAddress, acquire_delay: float):
  25. """
  26. Constructor.
  27. :param port: serial port
  28. :param module_address: address of module
  29. :param device_cnt: number of devices
  30. """
  31. self.__port = port
  32. self.__module_addr = module_address
  33. self.__params = {}
  34. self.__devices = {}
  35. self._lock = Lock() # thread lock
  36. self.__in_acquisition = False # flag when in acquisition
  37. self.__acquire_delay = acquire_delay # intermediate delay between data fetching when data not ready
  38. self.__acquire_thread = None # thread for data acquisition
  39. self.__data_raw = bytearray() # raw data buffer
  40. self.__df = pd.DataFrame() # internal data frame
  41. self.__ts0 = None # 0-th timestamp
  42. self.__payload_size = 0
  43. """
  44. ================================================================================================================
  45. Module General Methods
  46. ================================================================================================================
  47. """
  48. def init(self) -> None:
  49. """
  50. Initialize the module from hardware.
  51. """
  52. sensors = self.cmd_get_sensor_type()
  53. for sen_address, sen_type in sensors.items():
  54. self.__devices[sen_address] = NBusSensor(self.__port, self.__module_addr, sen_address)
  55. self.__devices[sen_address].type = sen_type
  56. self.cmd_get_format()
  57. self.__calculate_payload_size()
  58. def get_devices(self) -> dict[NBusSensorAddress, NBusSensor]:
  59. """
  60. Get module devices.
  61. :return: dictionary of connected devices
  62. """
  63. return self.__devices
  64. """
  65. ================================================================================================================
  66. Module Get Commands
  67. ================================================================================================================
  68. """
  69. def cmd_get_echo(self, message: bytearray) -> bool:
  70. """
  71. Get echo from module.
  72. :param message: message to send
  73. :return: status (True = echo, False = no echo)
  74. """
  75. _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_ECHO, message)
  76. return response == list(message)
  77. def cmd_get_param(self, parameter: NBusParameterID) -> NBusParameterValue:
  78. """
  79. Get single module parameter.
  80. :param parameter: parameter id
  81. :return: parameter value
  82. """
  83. # get response
  84. resp_len, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_PARAM,
  85. bytearray([parameter.value]))
  86. # parse parameter
  87. param_id, param_val = NbusCommonParser.parameters_from_response(resp_len, response)[0]
  88. # store parameter value
  89. self.__params[param_id] = param_val
  90. return param_val
  91. def cmd_get_all_params(self) -> dict[NBusParameterID, NBusParameterValue]:
  92. """
  93. Get all module parameters.
  94. :return: dict of parameter id and parameter value
  95. """
  96. # get response
  97. resp_len, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_PARAM, bytearray([]))
  98. # parse parameters
  99. params = NbusCommonParser.parameters_from_response(resp_len, response)
  100. for param_id, param_val in params:
  101. # store parameters
  102. self.__params[param_id] = param_val
  103. return self.__params.copy()
  104. def cmd_get_sensor_cnt(self) -> NBusSensorCount:
  105. """
  106. Get sensor count.
  107. :return: count of read-only and read-write sensors
  108. """
  109. _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_SENSOR_CNT, bytearray([]))
  110. return NBusSensorCount(*response)
  111. def cmd_get_data(self) -> dict[NBusSensorAddress, list[NBusDataValue]]:
  112. """
  113. Get data from all module sensors.
  114. :return: dict of device addresses and data values
  115. """
  116. # get data
  117. resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_DATA, bytearray([]))
  118. # parse data
  119. begin_idx = 0
  120. data = {}
  121. while begin_idx < resp_length:
  122. device_id = response[begin_idx]
  123. # handle errors
  124. if self.__devices[device_id].data_format is None: # check for format and params
  125. raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED)
  126. values, offset = NbusCommonParser.data_from_response(self.__devices[device_id].data_format,
  127. response[begin_idx:])
  128. data[device_id] = values
  129. begin_idx += offset + 1
  130. return data
  131. def cmd_get_info(self) -> NBusModuleInfo:
  132. """
  133. Get module info.
  134. :return: module info
  135. """
  136. response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_INFO, bytearray([]))
  137. name = str(response[1:9], "ascii")
  138. typ = str(response[9:12], "ascii")
  139. uuid = struct.unpack("<I", bytearray(response[12:16]))[0]
  140. hw = str(response[16:19], "ascii")
  141. fw = str(response[19:22], "ascii")
  142. mem_id = struct.unpack("<Q", bytearray(response[22:30]))[0]
  143. ro_count = int(response[30])
  144. rw_count = int(response[31])
  145. return NBusModuleInfo(module_name=name, module_type=typ, uuid=uuid, hw=hw, fw=fw, memory_id=mem_id,
  146. read_only_sensors=ro_count, read_write_sensors=rw_count)
  147. def cmd_get_format(self) -> dict[NBusSensorAddress, NBusDataFormat]:
  148. """
  149. Get format of all on-board sensors.
  150. :return: dict of sensor addresses and sensor formats
  151. """
  152. # get response
  153. resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_FORMAT,
  154. bytearray([]))
  155. begin_idx = 0
  156. formats = {}
  157. # parse format
  158. while begin_idx < resp_length:
  159. device_id = response[begin_idx]
  160. device_format = NbusCommonParser.format_from_response(response[begin_idx:begin_idx + 4])
  161. self.__devices[device_id].data_format = device_format
  162. formats[device_id] = device_format
  163. begin_idx += 4
  164. return formats
  165. def cmd_get_sensor_type(self) -> dict[NBusSensorAddress, NBusSensorType]:
  166. """
  167. Get type of all on-board sensors.
  168. :return: dict of sensor addresses and sensor types
  169. """
  170. # get response
  171. resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_GET_SENSOR_TYPE,
  172. bytearray([]))
  173. # parse response
  174. types = {}
  175. i = 0
  176. while i < resp_length - 1:
  177. types[NBusSensorAddress(response[i])] = NBusSensorType(response[i + 1])
  178. i += 2
  179. return types
  180. """
  181. ================================================================================================================
  182. Module Set Commands
  183. ================================================================================================================
  184. """
  185. def cmd_set_module_stop(self) -> None:
  186. """
  187. Stop automatic measuring.
  188. :return: status
  189. """
  190. self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_STOP, bytearray([]))
  191. def cmd_set_module_start(self) -> None:
  192. """
  193. Start automatic measuring.
  194. :return: status
  195. """
  196. self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_START, bytearray([]))
  197. def cmd_set_param(self, param: NBusParameterID, value: NBusParameterValue) -> NBusStatusType:
  198. """
  199. Set module parameter.
  200. :param param: parameter ID
  201. :param value: parameter value
  202. :return: status
  203. """
  204. # create request packet
  205. param_id_raw = struct.pack("B", param.value)
  206. param_val_raw = struct.pack("<I",value)
  207. param_bytes = bytearray(param_id_raw) + bytearray(param_val_raw)
  208. # proceed request
  209. _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_PARAM, param_bytes)
  210. # if response is valid, store parameter
  211. if response[0] == param.value and response[1] == NBusStatusType.STATUS_SUCCESS:
  212. self.__params[param] = value
  213. return NBusStatusType(response[1])
  214. def cmd_set_multi_params(self, params: dict[NBusParameterID, NBusParameterValue]) \
  215. -> dict[NBusParameterID, NBusStatusType]:
  216. """
  217. Set multiple module parameters.
  218. :param params: parameters
  219. :return: dict od statuses
  220. """
  221. # create request packet
  222. param_bytes = NbusCommonParser.parameters_to_request(params)
  223. # send request
  224. resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_PARAM,
  225. bytearray(param_bytes), long_answer=1.0)
  226. # parse statuses
  227. statuses = {}
  228. for i in range(0, resp_length - 1, 2):
  229. p_id = NBusParameterID(response[i])
  230. if response[i+1] == NBusStatusType.STATUS_SUCCESS:
  231. self.__params[p_id] = params[p_id] # if success, store param
  232. statuses[p_id] = NBusStatusType(response[i + 1])
  233. return statuses
  234. def cmd_set_calibrate(self) -> NBusStatusType:
  235. """
  236. Send calibration command.
  237. :return: calibration status
  238. """
  239. resp_length, *response = self.__port.request_module(self.__module_addr,
  240. NBusCommand.CMD_SET_CALIBRATE, bytearray([]))
  241. print(response)
  242. return NBusStatusType(response[0])
  243. def cmd_set_data(self, data: dict[NBusSensorAddress, list[NBusDataValue]]) \
  244. -> dict[NBusSensorAddress, NBusStatusType]:
  245. """
  246. Set data to read-write sensors.
  247. :param data:
  248. :return: operation statuses
  249. """
  250. # create request packet
  251. request = []
  252. # transform data
  253. for addr in data.keys():
  254. # handle errors
  255. if self.__devices[addr].data_format is None: # check for format and params
  256. raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED)
  257. raw_data = data[addr]
  258. request.append(addr)
  259. request.extend(NbusCommonParser.data_to_request(self.__devices[addr].data_format, raw_data))
  260. # send request
  261. resp_length, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_DATA,
  262. bytearray(request))
  263. # return response statuses
  264. statuses = {}
  265. for i in range(0, resp_length - 1, 2):
  266. statuses[NBusSensorAddress(response[i])] = NBusStatusType(response[i + 1])
  267. return statuses
  268. """
  269. ================================================================================================================
  270. Automatic Data Stream Methods
  271. ================================================================================================================
  272. """
  273. def start_streaming(self) -> None:
  274. """
  275. Start data streaming (e.g. auto-cast).
  276. """
  277. self.cmd_set_module_start()
  278. self.__acquire_thread = Thread(target=self.__acquire_callback)
  279. self.__in_acquisition = True
  280. # end thread if running
  281. if self.__acquire_thread is not None and self.__acquire_thread.is_alive():
  282. self.__acquire_thread.join()
  283. self.__acquire_thread.start()
  284. def stop_streaming(self):
  285. """
  286. Stop data streaming (e.g. auto-cast).
  287. """
  288. self.cmd_set_module_stop()
  289. self.__in_acquisition = False
  290. if self.__acquire_thread is not None and self.__acquire_thread.is_alive():
  291. self.__acquire_thread.join()
  292. self.__port.flush()
  293. def fetch_stream_chunk(self) -> pd.DataFrame:
  294. """
  295. Fetch data from stream (e.g. auto-cast).
  296. Can be called anytime.
  297. It not erase internal dataframe.
  298. :return: stream data frame
  299. """
  300. with self._lock:
  301. packets = self.__data_raw.split(NBUS_BRIDGE_DATA_HDR)[1:]
  302. packet_cnt = len(packets) - self.__in_acquisition
  303. parsed_packets = []
  304. # parse packets
  305. for i in range(packet_cnt):
  306. data = self._parse_data_from_stream_packet(packets[i])
  307. if data is not None:
  308. parsed_packets.append(data)
  309. else:
  310. print("damaged: ", i, packets[i])
  311. # extend internal dataframe
  312. if parsed_packets:
  313. data_frame = pd.DataFrame(parsed_packets)
  314. self.__df = pd.concat([self.__df, data_frame], ignore_index=True).copy(deep=True)
  315. else:
  316. data_frame = pd.DataFrame()
  317. # erase raw data buffer
  318. if self.__in_acquisition:
  319. unparsed_bytes = len(packets[-1]) + NBUS_BRIDGE_DATA_HDR_SIZE
  320. self.__data_raw = self.__data_raw[-unparsed_bytes:]
  321. else:
  322. self.__data_raw = bytearray()
  323. self._transform_timestamp(data_frame)
  324. return data_frame
  325. def fetch_full_stream(self) -> pd.DataFrame:
  326. """
  327. Fetch all data from stream (e.g. auto-cast).
  328. Must be called after stop_streaming() method.
  329. It will erase internal dataframe.
  330. :return: stream dataframe
  331. """
  332. if self.__in_acquisition:
  333. return pd.DataFrame()
  334. self.fetch_stream_chunk()
  335. df = self.__df
  336. self._transform_timestamp(df)
  337. self.__df = pd.DataFrame()
  338. self.__ts0 = None
  339. return df
  340. """
  341. ================================================================================================================
  342. Internal Helper Methods
  343. ================================================================================================================
  344. """
  345. def _transform_timestamp(self, data_frame: pd.DataFrame) -> None:
  346. """
  347. Transform timestamp values in dataframe.
  348. :param data_frame: dataframe to transform
  349. """
  350. if not data_frame.empty and "TS" in self.__df.columns:
  351. if self.__ts0 is None:
  352. self.__ts0 = self.__df["TS"].iloc[0]
  353. data_frame["TS"] -= self.__ts0
  354. def _parse_data_from_stream_response(self, module_address: NBusModuleAddress, resp_length: int, response: bytearray) \
  355. -> dict[str, NBusDataValue]:
  356. """
  357. Parse data of slave from stream response.
  358. :param module_address: address of module
  359. :param resp_length: length of response
  360. :param response: raw data
  361. :return: dict of data values
  362. """
  363. data_offset = 0
  364. data = {}
  365. while data_offset < resp_length:
  366. device_id = response[data_offset]
  367. if self.__devices[device_id].data_format is None:
  368. raise NBusErrorAPI(NBusErrorAPIType.FORMAT_NOT_LOADED)
  369. values, offset = NbusCommonParser.data_from_response(self.__devices[device_id].data_format, response[data_offset:])
  370. data_tag = str(module_address) + "." + str(device_id)
  371. for i in range(len(values)):
  372. data[data_tag + "." + str(i + 1)] = values[i]
  373. data_offset += offset + NBUS_SA_SIZE
  374. return data
  375. def _parse_data_from_stream_packet(self, data_packet: bytearray) \
  376. -> Optional[dict[str, NBusDataValue]]:
  377. """
  378. Parse the stream-mode data packet.
  379. :param data_packet: packet to parse
  380. :return: dictionary of data values or None
  381. """
  382. packet_size = len(data_packet)
  383. packet_crc = crc8(data_packet[:NBUS_CRC_ADDR])
  384. # check validity
  385. if packet_size < NBUS_TS_SIZE + NBUS_CRC_SIZE or data_packet[NBUS_CRC_ADDR] != packet_crc:
  386. return None
  387. # parse data
  388. try:
  389. data_offset = 0
  390. ts = struct.unpack("<I", data_packet[data_offset: data_offset + NBUS_TS_SIZE])[0]
  391. data = {"TS": ts}
  392. data_offset += NBUS_TS_SIZE
  393. module_addr = data_packet[data_offset]
  394. data_offset += NBUS_MA_SIZE
  395. packet = data_packet[data_offset: data_offset + self.__payload_size]
  396. data |= self._parse_data_from_stream_response(module_addr, self.__payload_size, packet)
  397. data_offset += self.__payload_size
  398. return data
  399. except Exception:
  400. return None
  401. def __calculate_payload_size(self):
  402. """
  403. Calculate payload size of stream-mode packet.
  404. :return: payload size in bytes
  405. """
  406. for device in self.__devices.values():
  407. fmt = device.data_format
  408. self.__payload_size += fmt.byte_length * fmt.samples + NBUS_SA_SIZE
  409. def __acquire_callback(self) -> None:
  410. """
  411. Thread worker for data acquisition.
  412. """
  413. while self.__in_acquisition:
  414. data = self.__port.try_read()
  415. if not data:
  416. time.sleep(self.__acquire_delay)
  417. continue
  418. with self._lock:
  419. self.__data_raw.extend(data)