tmpConnit

chiyu
Chiyu Chen 1 month ago
parent de8b153f92
commit 99e42dd996

@ -29,7 +29,7 @@ from .utils import acquireSerial, acquirePort
from .utils.acquirePort import find_available_port from .utils.acquirePort import find_available_port
logger = setup_logger(os.path.basename(__file__)) logger = setup_logger(os.path.basename(__file__))
PROJECT_VER = "v1.20" PROJECT_VER = "v1.25"
class PanelState: class PanelState:
def __init__(self): def __init__(self):
@ -1696,4 +1696,8 @@ if __name__ == "__main__":
2025.05.28 2025.05.28
1. 工程模式 新增 Node 重啟的功能 1. 工程模式 新增 Node 重啟的功能
2026.08.05 (ver 1.25)
1. 增加對於 esp32 連線模式的支援
''' '''

@ -277,11 +277,12 @@ class RFModule:
self.status = RFStatus() self.status = RFStatus()
self.socket_info = SocketInfo() self.socket_info = SocketInfo()
self.poll_data_times: deque[int] = deque([0] * N,maxlen=N) # ms # # only for ESP32
self.poll_data_sizes: deque[int] = deque([0] * N, maxlen=N) # num # self.poll_data_times: deque[int] = None
self.poll_total_time = 0 # ms # self.poll_data_sizes: deque[int] = None
self.poll_total_size = 0 # num # self.poll_total_time = None
self.remain_bytes = 0 # num # self.poll_total_size = None
# self.remain_bytes = None
def update_avg_bandwidth(self, bandwidth: float, timestamp: Optional[float] = None) -> None: def update_avg_bandwidth(self, bandwidth: float, timestamp: Optional[float] = None) -> None:
self.status.avg_bandwidth = bandwidth self.status.avg_bandwidth = bandwidth

@ -102,12 +102,42 @@ class FrameProcessor(ABC):
class RawFrameProcessor(FrameProcessor): class RawFrameProcessor(FrameProcessor):
"""原始數據直通處理器""" """原始數據直通處理器"""
def __init__(self):
N = 10
self.tx_interval_bytes: deque[int] = deque([0] * N,maxlen=N)
self.tx_interval_rate: deque[float] = deque([0] * N,maxlen=N)
self.tx_record_byte: int = 0
self.tx_record_time = time.time()
self.rx_interval_bytes: deque[int] = deque([0] * N, maxlen=N)
self.rx_interval_rate: deque[float] = deque([0] * N,maxlen=N)
self.rx_record_byte: int = 0
self.rx_record_time = time.time()
def process_incoming(self, data: bytes) -> bytes: def process_incoming(self, data: bytes) -> bytes:
"""直接返回原始數據,不進行緩衝""" """直接返回原始數據,不進行緩衝"""
interval = time.time() - self.tx_record_time
if interval >= 1:
self.tx_interval_bytes.append(self.tx_record_byte)
self.tx_interval_rate.append(self.tx_record_byte/interval)
self.tx_record_time = time.time()
self.tx_record_byte = 0
else:
self.tx_record_byte += len(data)
return [data] if data else [] return [data] if data else []
def process_outgoing(self, data: bytes) -> bytes: def process_outgoing(self, data: bytes) -> bytes:
"""直接返回原始數據,不進行封裝""" """直接返回原始數據,不進行封裝"""
interval = time.time() - self.rx_record_time
if interval >= 1:
self.rx_interval_bytes.append(self.rx_record_byte)
self.tx_interval_rate.append(self.tx_record_byte/interval)
self.rx_record_time = time.time()
self.rx_record_byte = 0
else:
self.rx_record_byte += len(data)
return data return data
@ -300,7 +330,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
MAX_BYTES_PER_FLUSH = 150 MAX_BYTES_PER_FLUSH = 150
MAX_PAYLOAD_PER_FRAME = 80 MAX_PAYLOAD_PER_FRAME = 80
CHUNK_SEND_INTERVAL_SEC = 0.005 CHUNK_SEND_INTERVAL_SEC = 0.020
@ -316,6 +346,10 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
self.received_len = 0 # 收到封包累計 self.received_len = 0 # 收到封包累計
self.poll_miss_count = 0 # poll 封包連續丟失的累計數 self.poll_miss_count = 0 # poll 封包連續丟失的累計數
self.esp_state = 1 # 紀錄每個 esp 連線狀態
# 1 : 物件剛創建 剛接到 HELO 但是還沒有在 mavlinkVehicleView 建立物件
# 2 : 已經擁有 mavlinkVehicleView 物件
# 9 : 連續的 poll failure 視為斷線
def __init__(self, at_handler: "ATCommandHandler" = None): def __init__(self, at_handler: "ATCommandHandler" = None):
super().__init__(at_handler) super().__init__(at_handler)
@ -586,17 +620,6 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
system_id, sender_address_64, time.time() system_id, sender_address_64, time.time()
) )
vehicle = vehicle_registry.get(system_id)
if vehicle and not vehicle.rf_module:
N = 10
vehicle.rf_module = RFModule(RFModuleType.XBEE_ESP)
vehicle.rf_module.poll_data_times: deque[int] = deque([0] * N,maxlen=N) # ms
vehicle.rf_module.poll_data_sizes: deque[int] = deque([0] * N, maxlen=N) # num
vehicle.rf_module.poll_total_time = 0 # ms
vehicle.rf_module.poll_total_size = 0 # num
vehicle.rf_module.remain_bytes = 0 # num
logger.debug(f"mark A")
logger.debug( logger.debug(
f"new HELO system_id={system_id}, address_64={sender_address_64.hex()}") f"new HELO system_id={system_id}, address_64={sender_address_64.hex()}")
else: else:
@ -818,16 +841,30 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
# 主要迴圈 Main Loop # 主要迴圈 Main Loop
async def _operator_tick(self) -> None: async def _operator_tick(self) -> None:
# 1. 更新資訊 # 1. 更新資訊 對於每個 ESP 裝置的資訊更新
# for src64, INFO in self.esp32_address_mapping.items(): for src64, INFO in self.esp32_address_mapping.items():
# vehicle = vehicle_registry.get(INFO.system_id)
# if vehicle and vehicle.rf_module:
# speed = vehicle.rf_module.poll_total_size / (vehicle.rf_module.poll_total_time + 1) vehicle = vehicle_registry.get(INFO.system_id)
# logger.debug(f"now speed: {speed:.2f} / total_size: {vehicle.rf_module.poll_total_size}") # 初始化 mavlinkVehicleView
if vehicle and INFO.esp_state == 1 :
N = 10
vehicle.rf_module = RFModule(RFModuleType.XBEE_ESP)
vehicle.rf_module.poll_data_times: deque[int] = deque([0] * N,maxlen=N) # ms
vehicle.rf_module.poll_data_sizes: deque[int] = deque([0] * N, maxlen=N) # num
vehicle.rf_module.poll_total_time = 0 # ms
vehicle.rf_module.poll_total_size = 0 # num
vehicle.rf_module.remain_bytes = 0 # num
vehicle.rf_module.rssi = 0 # db
INFO.esp_state = 2
elif vehicle and INFO.esp_state == 2:
pass
# speed = vehicle.rf_module.poll_total_size / (vehicle.rf_module.poll_total_time + 1)
# logger.debug(f"now speed: {speed:.2f} / total_size: {vehicle.rf_module.poll_total_size}")
# 2. 移除沒反應 dongle # 2. 移除沒反應 dongle
self._prune_stale_devices() self._prune_stale_devices()
@ -1591,25 +1628,30 @@ if __name__ == '__main__':
serial_id = 1 serial_id = 1
processor = sm.get_espv1_processor(serial_id) processor = sm.get_espv1_processor(serial_id)
if processor is not None: if processor is None:
# for x in range(10): return
# processor.request_discovery()
# time.sleep(1) # for x in range(10):
# print(f"mark N{x}") # processor.request_discovery()
# processor.request_poll(target_system_id=1) # time.sleep(1)
# processor.request_poll(target_system_id=1, grant_bytes=200) # print(f"mark N{x}")
print(processor.get_status_snapshot()) # processor.request_poll(target_system_id=1)
print(processor.get_gcs_queue_byte_count()) # processor.request_poll(target_system_id=1, grant_bytes=200)
print(processor.get_status_snapshot())
print(processor.get_gcs_queue_byte_count())
# rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0xFE) # rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0xFE)
# for i in range(3): # for i in range(3):
# processor = sm.get_espv1_processor(serial_id) # deadline = time.time() + 5.0 # 最多等 5 秒,避免永遠卡住
# if processor is not None: # while processor.operator_busy and time.time() < deadline:
# deadline = time.time() + 5.0 # 最多等 5 秒,避免永遠卡住 # time.sleep(0.01)
# while processor.operator_busy and time.time() < deadline: # sm.send_at_command(1, rssi_request)
# time.sleep(0.01) # time.sleep(1)
# sm.send_at_command(1, rssi_request)
# time.sleep(1) for i in range(20):
print(f"received: {processor.tx_interval_bytes}")
time.sleep(1)
time.sleep(120) time.sleep(120)
sm.remove_serial_link(serial_id) sm.remove_serial_link(serial_id)

Loading…
Cancel
Save