|
|
@@ -1,9 +1,16 @@
|
|
|
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
|
|
|
@@ -18,7 +25,7 @@ class NBusSlaveModule:
|
|
|
Class representing nBus slave module.
|
|
|
"""
|
|
|
|
|
|
- def __init__(self, port: NBusPort, module_address: NBusModuleAddress):
|
|
|
+ def __init__(self, port: NBusPort, module_address: NBusModuleAddress, acquire_delay: float):
|
|
|
"""
|
|
|
Constructor.
|
|
|
|
|
|
@@ -31,16 +38,24 @@ class NBusSlaveModule:
|
|
|
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, load_format: bool) -> None:
|
|
|
+ def init(self) -> None:
|
|
|
"""
|
|
|
Initialize the module from hardware.
|
|
|
-
|
|
|
- :param load_format: flag to fetch data format
|
|
|
"""
|
|
|
sensors = self.cmd_get_sensor_type()
|
|
|
|
|
|
@@ -48,8 +63,8 @@ class NBusSlaveModule:
|
|
|
self.__devices[sen_address] = NBusSensor(self.__port, self.__module_addr, sen_address)
|
|
|
self.__devices[sen_address].type = sen_type
|
|
|
|
|
|
- if load_format:
|
|
|
- self.cmd_get_format()
|
|
|
+ self.cmd_get_format()
|
|
|
+ self.__calculate_payload_size()
|
|
|
|
|
|
def get_devices(self) -> dict[NBusSensorAddress, NBusSensor]:
|
|
|
"""
|
|
|
@@ -218,24 +233,21 @@ class NBusSlaveModule:
|
|
|
Module Set Commands
|
|
|
================================================================================================================
|
|
|
"""
|
|
|
- def cmd_set_module_stop(self) -> NBusStatusType:
|
|
|
+ def cmd_set_module_stop(self) -> None:
|
|
|
"""
|
|
|
Stop automatic measuring.
|
|
|
|
|
|
:return: status
|
|
|
"""
|
|
|
+ self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_STOP, bytearray([]))
|
|
|
|
|
|
- _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_STOP, bytearray([]))
|
|
|
- return NBusStatusType(response[0])
|
|
|
-
|
|
|
- def cmd_set_module_start(self) -> NBusStatusType:
|
|
|
+ def cmd_set_module_start(self) -> None:
|
|
|
"""
|
|
|
Start automatic measuring.
|
|
|
|
|
|
:return: status
|
|
|
"""
|
|
|
- _, *response = self.__port.request_module(self.__module_addr, NBusCommand.CMD_SET_START, bytearray([]))
|
|
|
- return NBusStatusType(response[0])
|
|
|
+ self.__port.send_module(self.__module_addr, NBusCommand.CMD_SET_START, bytearray([]))
|
|
|
|
|
|
def cmd_set_param(self, param: NBusParameterID, value: NBusParameterValue) -> NBusStatusType:
|
|
|
"""
|
|
|
@@ -332,3 +344,200 @@ class NBusSlaveModule:
|
|
|
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)
|