From fd949fca47524d8420417a19265c3e87b468c63f Mon Sep 17 00:00:00 2001 From: Chiyu Chen Date: Sat, 27 Jun 2026 05:50:39 +0800 Subject: [PATCH] (modify) serialManager.py MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 修正 POLL-DONE 程序的 timeout 判定方式 - 修正為 py3.8 語法 - 不在判定要求量跟實收量一致性 --- .../fc_network_adapter/mavlinkObject.py | 4 +- .../fc_network_adapter/mavlinkVehicleView.py | 2 + .../fc_network_adapter/serialManager.py | 218 +++++++++++------- .../fc_network_adapter/utils/pollStrategy.py | 23 -- 4 files changed, 139 insertions(+), 108 deletions(-) diff --git a/src/fc_network_adapter/fc_network_adapter/mavlinkObject.py b/src/fc_network_adapter/fc_network_adapter/mavlinkObject.py index f1ab903..9e94b10 100644 --- a/src/fc_network_adapter/fc_network_adapter/mavlinkObject.py +++ b/src/fc_network_adapter/fc_network_adapter/mavlinkObject.py @@ -47,8 +47,8 @@ from pymavlink.dialects.v20 import ardupilotmega as mav_ardupilot # 自定義的 import from .mavlinkVehicleView import ( - vehicle_registry, - VehicleView, + vehicle_registry, # 儲存全部物件的地方 + VehicleView, # 代表每台載具的最外層 VehicleComponent, ComponentType, StatusTextEntry, diff --git a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py index 89b3cb4..7ccef8d 100644 --- a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py +++ b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py @@ -187,6 +187,7 @@ class RFStatus: @dataclass class SocketInfo: """Socket連接資訊""" + src64_addr: Optional[bytes] = None # 模組的物理定址 ip: Optional[str] = None # IP位址 port: Optional[int] = None # 埠號 local_ip: Optional[str] = None # 本地IP @@ -310,6 +311,7 @@ class RFModule: class VehicleView: """ + 最上層 載具視圖 - 純狀態容器 特點: diff --git a/src/fc_network_adapter/fc_network_adapter/serialManager.py b/src/fc_network_adapter/fc_network_adapter/serialManager.py index 8e321b7..53714e8 100644 --- a/src/fc_network_adapter/fc_network_adapter/serialManager.py +++ b/src/fc_network_adapter/fc_network_adapter/serialManager.py @@ -25,6 +25,12 @@ from typing import List # from xbee.frame import APIFrame # 自定義的 import +from .mavlinkVehicleView import ( + vehicle_registry, # 儲存全部物件的地方 + VehicleView, + RFModule, + RFModuleType, +) from .utils import RingBuffer, setup_logger from .utils import pollStrategy @@ -139,6 +145,8 @@ class XBeeFrameProcessor_Base(FrameProcessor): DEST_ADDR16_UNICAST = b'\xFF\xFE' DEST_ADDR16_BRAODCAST = b'\xFF\xFF' + DEST_ADDR64_BRAODCAST = b'\x00\x00\x00\x00\x00\x00\x00\x00' + def __init__(self, at_handler: "ATCommandHandler" = None): super().__init__() self.at_handler = at_handler @@ -222,7 +230,7 @@ class XBeeFrameProcessor_Base(FrameProcessor): @staticmethod def _encapsulate( data: bytes, - dest_addr64: bytes = b'\x00\x00\x00\x00\x00\x00\x00\x00', + dest_addr64: bytes = DEST_ADDR64_BRAODCAST, dest_addr16 = DEST_ADDR16_BRAODCAST, frame_id: int = 0x01, ) -> bytes: @@ -252,14 +260,26 @@ class XBeeFrameProcessor_Base(FrameProcessor): class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): - - # GCS -> UAV: - # DISC - # POLL + esp_sysid(1) + grant_bytes(2) - # - # UAV -> GCS: - # HELO + esp_sysid(1) - # DONE + sysid(1) + sent_len(2) + remain_len(2) + ''' + ESP32 封包分類 + - Mavlink 資料封包 : + - Payload = mavlink [方向] GCS -> UAV (broadcast) + - Payload = mavlink [方向] UAV -> GCS (with src64 Addr) + - Discovery 階段 + - Payload = DISC_HEADER [方向] GCS -> UAV (broadcast) + - Payload = HELLO_HEADER + esp_sysid(1) [方向] UAV -> GCS (with src64 Addr) + - Poll 階段 + - Payload = POLL_HEADER + esp_sysid(1) + grant_bytes(2) [方向] GCS -> UAV (with src64 Addr) + - Payload = DONE_HEADER + sysid(1) + sent_len(2) + remain_len(2) [方向] UAV -> GCS (with src64 Addr) + + 這邊認定 esp_sysid 會等於 mavlink sysid + + 硬體產品系列 (Product Family) : XBP9B-DM + 晶片世代: XBee PRO 900HP 200K + 運作頻段: 900 MHz RF。 + 運行協議 (Protocol / Function Set): DigiMesh + 韌體版本 (Firmware Version): 8075 + ''' DISC_HEADER = b'DISC' HELLO_HEADER = b'HELO' @@ -270,6 +290,11 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): MAX_PAYLOAD_PER_FRAME = 80 CHUNK_SEND_INTERVAL_SEC = 0.01 + # ADDR16 選項 + DEST_ADDR16_BRAODCAST = b'\xFF\xFF' + + DEST_ADDR64_BRAODCAST = b'\x00\x00\x00\x00\x00\x00\xFF\xFF' + class Esp32DeviceInfo: def __init__(self, system_id, address_64, last_hello_time): self.system_id = system_id @@ -277,6 +302,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.last_hello_time = last_hello_time self.remain_bytes = 0 # 剩餘 buffer 量 + self.last_poll_time = 0.0 self.last_done_time = 0.0 # 最後送出Done的時間 self.received_len = 0 # 收到封包累計 @@ -301,7 +327,9 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.poll_done_event: Optional[asyncio.Event] = None self.current_poll_address_64: Optional[bytes] = None - self.last_discovery_time = 0.0 + self.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間 + self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的 + self.MAX_mavPack_interval_timeout = 100 # ms,poll 期間 MAVLink/DONE 最大閒置間隔 self.discovery_interval_seconds = 30.0 # 每次做 discovery 程序的間隔時間 self.device_offline_timeout = self.discovery_interval_seconds * 2 # 遠端沒有回應會被踢出 超時時限 self.operator_tick_interval_seconds = 0.03 # @@ -352,6 +380,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): remote_device = self.esp32_address_mapping.get(sender_address_64) if remote_device is not None: remote_device.received_len += len(payload) + if (self.operator_busy and self.current_poll_address_64 == sender_address_64): + self.last_recieve_mavlink = time.time() return payload if frame_type == self.FRAME_TYPE_AT_RESPONSE: @@ -376,7 +406,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): def pack_discovery(self) -> bytes: return self._encapsulate(self.DISC_HEADER, frame_id=0x00) - def handle_hello_report(self, payload: bytes, sender_address_64: bytes) -> None: + # 處理每個裝置回傳的 Hello 訊息 + def handle_hello_report(self, payload: bytes, sender_address_64: bytes) -> None: system_id = payload[4] remote_device = self.esp32_address_mapping.get(sender_address_64) @@ -385,8 +416,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): system_id, sender_address_64, time.time() ) 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()}") elif remote_device.address_64 == sender_address_64: remote_device.last_hello_time = time.time() else: @@ -394,6 +424,10 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): f"SYSID duplicated system_id={system_id}, " f"address_64={sender_address_64.hex()}" ) + + if vehicle:=vehicle_registry.get(system_id): + if not vehicle.rf_module: + vehicle.rf_module = RFModule(RFModuleType.XBEE) def pack_poll(self, target_address_64: bytes, grant_bytes: int = 0) -> Optional[bytes]: remote_device = self.esp32_address_mapping.get(target_address_64) @@ -416,13 +450,16 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): if remote_device is None: return - system_id, sent_length, remain_length = struct.unpack('>BHH', payload[4:9]) - if sent_length != remote_device.received_len: - logger.info( - f"POLL may be missing packets sent={sent_length} " - f"received={remote_device.received_len} system_id={system_id}" - ) + # 這段是有問題的 因為會有整數封包切割問題 以及載具端的 buffer 存量不足 故回傳的資訊量會與要求的不一致 + # system_id, sent_length, remain_length = struct.unpack('>BHH', payload[4:9]) + # if sent_length != remote_device.received_len: + # logger.info( + # f"POLL may be missing packets sent={sent_length} " + # f"received={remote_device.received_len} system_id={system_id}" + # ) + # TODO 傳送速率 + # TODO 累積速率預測 remote_device.received_len = 0 remote_device.remain_bytes = remain_length remote_device.last_done_time = time.time() @@ -716,6 +753,22 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.is_discovery_phase = False self.pending_manual_discovery = False + async def _wait_poll_done_with_idle_timeout(self) -> bool: + idle_timeout_sec = self.MAX_mavPack_interval_timeout / 1000.0 + poll_tick = min(0.02, idle_timeout_sec / 2) + + while not self.poll_done_event.is_set(): + if time.time() - self.last_recieve_mavlink >= idle_timeout_sec: + return False + try: + await asyncio.wait_for( + self.poll_done_event.wait(), + timeout=poll_tick, + ) + except asyncio.TimeoutError: + continue + return True + # poll 程序 async def _run_one_poll(self, target_address_64: bytes, grant_bytes: int) -> None: if self.serial_writer is None: @@ -734,22 +787,15 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.current_poll_address_64 = target_address_64 self.poll_done_event.clear() self.serial_writer(poll_frame) + self.last_recieve_mavlink = time.time() - timeout_seconds = pollStrategy.estimate_poll_timeout( - grant_bytes, - self.serial_baudrate, - self.guard_milliseconds, - ) try: - await asyncio.wait_for( - self.poll_done_event.wait(), - timeout=timeout_seconds, - ) - except asyncio.TimeoutError: - logger.warning( - f"POLL timeout address_64={target_address_64.hex()} " - f"grant_bytes={grant_bytes}" - ) + completed = await self._wait_poll_done_with_idle_timeout() + if not completed: + logger.warning( + f"POLL timeout address_64={target_address_64.hex()} " + f"grant_bytes={grant_bytes}" + ) finally: self.operator_busy = False self.current_poll_address_64 = None @@ -820,21 +866,21 @@ class ATCommandHandler: # ---- 接收端 ---- def handle_frame(self, frame: bytes) -> None: """ - 接收一整個 AT Response frame: + 接收一整個 AT Response frame: 1. 解析成 ATResponse 2. 推進 rx_module_ack 供其他模組消費 3. 本地 dispatch 給對應的 _handle_xxx """ - parsed = self._parse(frame) - if parsed is None: + parsed_at_ack = self._parse(frame) + if parsed_at_ack is None: return - if not rx_module_ack.put(parsed): + if not rx_module_ack.put(parsed_at_ack): logger.warning( - f"[{self.serial_port}] rx_module_ack overflow, drop {parsed.command!r}" + f"[{self.serial_port}] rx_module_ack overflow, drop {parsed_at_ack.command!r}" ) - self._dispatch(parsed) + self._dispatch(parsed_at_ack) @staticmethod def _parse(frame: bytes) -> ATResponse: @@ -865,26 +911,27 @@ class ATCommandHandler: handler = self.handlers.get(response.command) if handler: - handler(response.data) + handler(response) else: logger.debug( f"[{self.serial_port}] 未處理的 AT 指令: " f"{response.command.decode()}" ) - def _handle_rssi(self, data: bytes): + def _handle_rssi(self, response: ATResponse): """處理 DB (RSSI) 回應:單 byte 無號值,單位 dBm""" - pass - if data: - print(f"[{self.serial_port}] RSSI = -{data[0]} dBm") # dev - # logger.debug(f"[{self.serial_port}] RSSI = -{data[0]} dBm") # dev - pass - - def _handle_serial_high(self, data: bytes): + + # print(f"[{self.serial_port}] RSSI = -{data[0]} dBm") # dev + logger.debug(f"[{self.serial_port}] RSSI = -{response.data[0]} dBm") # dev + + + + + def _handle_serial_high(self, response: ATResponse): """處理 SH (Serial Number High)""" pass - def _handle_serial_low(self, data: bytes): + def _handle_serial_low(self, response: ATResponse): """處理 SL (Serial Number Low)""" pass @@ -1246,7 +1293,7 @@ class serial_manager: def send_at_command(self, serial_id, request: ATRequest) -> bool: """ - 對指定 serial_id 的 XBee dongle 發送一筆 AT 指令(thread-safe) + 對指定 serial_id 的 XBee dongle 發送一筆 AT 指令 (thread-safe) - serial_id: create_serial_link 取得的編號 - request: ATRequest 物件,攜帶 command / parameter / frame_id 回傳是否成功排進事件圈 @@ -1327,46 +1374,51 @@ if __name__ == '__main__': # UDP_REMOTE_PORT = 14571 # sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.STRAIGHT) - # # 測試項二 - # SERIAL_PORT = '/dev/ttyUSB0' # 手動指定 - # SERIAL_BAUDRATE = 115200 - # UDP_REMOTE_PORT = 14561 - # sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.XBEEAPI2AT) + # 測試項二 + print("運行 測試項二") + SERIAL_PORT = '/dev/ttyUSB0' # 手動指定 + SERIAL_BAUDRATE = 115200 + UDP_REMOTE_PORT = 14561 + sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.XBEEAPI2AT) + serial_id = 1 + device_sys_id = 10 - # linked_serial = sm.get_serial_link() - # print(linked_serial) + linked_serial = sm.get_serial_link() + print(f"連結完成 : {linked_serial}. 等待兩秒") - # # 等 connection_made 完成 writer 注入,再發一筆 AT 指令測試 - # time.sleep(5) - # rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0x52) - # for i in range(60): - # sm.send_at_command(1, rssi_request) - # time.sleep(1) + # 等 connection_made 完成 writer 注入,再發一筆 AT 指令測試 + time.sleep(2) + rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=device_sys_id) + print(f"手動送出 DB AT Command:") + for i in range(20): + sm.send_at_command(1, rssi_request) + time.sleep(1) - # sm.remove_serial_link(1) - # time.sleep(3) - # sm.shutdown() + sm.remove_serial_link(1) + time.sleep(2) + sm.shutdown() + print("結束運行") # # 測試項三 - SERIAL_PORT = '/dev/ttyUSB0' - SERIAL_BAUDRATE = 115200 - UDP_REMOTE_PORT = 14561 - sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.XBEEAPI_espv1) + # SERIAL_PORT = '/dev/ttyUSB0' + # SERIAL_BAUDRATE = 115200 + # UDP_REMOTE_PORT = 14561 + # sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.XBEEAPI_espv1) - time.sleep(2) # 等 serial 連線與 operator 啟動 + # time.sleep(2) # 等 serial 連線與 operator 啟動 - serial_id = 1 - processor = sm.get_espv1_processor(serial_id) - if processor is not None: - processor.request_discovery() - processor.request_poll(target_system_id=1) - processor.request_poll(target_system_id=1, grant_bytes=200) - print(processor.get_status_snapshot()) - print(processor.get_gcs_queue_byte_count()) - - sm.remove_serial_link(serial_id) - time.sleep(3) - sm.shutdown() + # serial_id = 1 + # processor = sm.get_espv1_processor(serial_id) + # if processor is not None: + # processor.request_discovery() + # processor.request_poll(target_system_id=1) + # processor.request_poll(target_system_id=1, grant_bytes=200) + # print(processor.get_status_snapshot()) + # print(processor.get_gcs_queue_byte_count()) + + # sm.remove_serial_link(serial_id) + # time.sleep(30) + # sm.shutdown() ''' ================= 改版記錄 ============================ diff --git a/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py b/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py index ca40ded..292c5d0 100644 --- a/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py +++ b/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py @@ -46,26 +46,3 @@ def pick_next( return selected_device.address_64, grant_bytes - -def estimate_poll_timeout( - grant_bytes: int, - baudrate: int, - guard_milliseconds: int = 20, -) -> float: - """ - 估算 POLL 後等待 DONE 的超時秒數。 - 移植自 udptest8 estimate_tdma_timeout。 - """ - grant_bytes = max(0, int(grant_bytes)) - max_payload_per_chunk = 100 - chunk_count = max(1, (grant_bytes + max_payload_per_chunk - 1) // max_payload_per_chunk) - - uart_time = (grant_bytes + chunk_count * 18 + 32) * 10.0 / baudrate - timeout = ( - uart_time - + chunk_count * 0.012 - + (guard_milliseconds / 1000.0) - + 0.35 - ) - return timeout -