From 1e07d1fa5438c1ecdd9fa9c48f34421ebc2e9140 Mon Sep 17 00:00:00 2001 From: Chiyu Chen Date: Mon, 20 Jul 2026 14:27:42 +0800 Subject: [PATCH] TempCommit --- .../fc_network_adapter/serialManager.py | 412 +++++++++++------- 1 file changed, 259 insertions(+), 153 deletions(-) diff --git a/src/fc_network_adapter/fc_network_adapter/serialManager.py b/src/fc_network_adapter/fc_network_adapter/serialManager.py index b874142..0d38e44 100644 --- a/src/fc_network_adapter/fc_network_adapter/serialManager.py +++ b/src/fc_network_adapter/fc_network_adapter/serialManager.py @@ -135,24 +135,25 @@ class XBeeFrameProcessor_Base(FrameProcessor): 韌體版本 (Firmware Version): 2014 -> 2 (XBee 3 平台) + 0 (802.15.4 協議代碼) + 14 (次要版本號)。 """ - # XBee API frame type - FRAME_TYPE_TX_REQUEST = 0x10 - FRAME_TYPE_AT_RESPONSE = 0x88 - FRAME_TYPE_TX_STATUS = 0x8B - FRAME_TYPE_RX_PACKET = 0x90 + # # XBee API frame type + # FRAME_TYPE_TX_REQUEST = 0x10 + # FRAME_TYPE_AT_RESPONSE = 0x88 + # FRAME_TYPE_TX_STATUS = 0x8B + # FRAME_TYPE_RX_PACKET = 0x90 - # ADDR16 選項 - DEST_ADDR16_UNICAST = b'\xFF\xFE' - DEST_ADDR16_BRAODCAST = b'\xFF\xFF' + # # ADDR16 選項 + # DEST_ADDR16_UNICAST = b'\xFF\xFE' + # DEST_ADDR16_BRAODCAST = b'\xFF\xFF' - DEST_ADDR64_BRAODCAST = b'\x00\x00\x00\x00\x00\x00\x00\x00' + # 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 - - + self.DEST_ADDR16_UNICAST = b'\xFF\xFE' + self.DEST_ADDR16_BRAODCAST = b'\xFF\xFF' + self.DEST_ADDR64_BRAODCAST = b'\x00\x00\x00\x00\x00\x00\x00\x00' # ---- 對外契約 ---- def process_incoming(self, data: bytes) -> bytes: @@ -172,7 +173,11 @@ class XBeeFrameProcessor_Base(FrameProcessor): def process_outgoing(self, data: bytes) -> bytes: """將數據封裝為 XBee API 傳輸幀""" - return self._encapsulate(data) + return self._encapsulate( + data, + dest_addr64 = self.DEST_ADDR64_BRAODCAST, + dest_addr16 = self.DEST_ADDR16_BRAODCAST, + frame_id = 0x01) # ---- 內部:拆幀與分派 ---- def _try_extract_frame(self) -> bytes: @@ -209,15 +214,15 @@ class XBeeFrameProcessor_Base(FrameProcessor): """ frame_type = frame[3] - if frame_type == self.FRAME_TYPE_RX_PACKET: # mavlink + if frame_type == ATCommandHandler.FRAME_TYPE_RX_PACKET: # mavlink return self._decapsulate(frame)[0] - if frame_type == self.FRAME_TYPE_AT_RESPONSE: # AT command + if frame_type == ATCommandHandler.FRAME_TYPE_AT_RESPONSE: # AT command if self.at_handler is not None: self.at_handler.handle_frame(frame) return None - if frame_type == self.FRAME_TYPE_TX_STATUS: + if frame_type == ATCommandHandler.FRAME_TYPE_TX_STATUS: # length = (frame[1] << 8) | frame[2] # for debug # logger.info( # f"TX Status raw={frame.hex()}, api_len={length}, " @@ -233,18 +238,18 @@ class XBeeFrameProcessor_Base(FrameProcessor): @staticmethod def _encapsulate( data: bytes, - dest_addr64: bytes = DEST_ADDR64_BRAODCAST, - dest_addr16: bytes = DEST_ADDR16_BRAODCAST, - frame_id: int = 0x01, + dest_addr64: bytes, + dest_addr16: bytes, + frame_id: int = 0x00, + options: int = 0x00, ) -> bytes: """ 將 payload 包成 XBee TX Request (0x10) - 使用廣播地址 - 添加適當的頭部和校驗和 """ - frame_type = XBeeFrameProcessor_Base.FRAME_TYPE_TX_REQUEST + frame_type = ATCommandHandler.FRAME_TYPE_TX_REQUEST broadcast_radius = 0x00 - options = 0x00 frame = struct.pack(">B", frame_type) + struct.pack(">B", frame_id) frame += dest_addr64 + dest_addr16 @@ -279,6 +284,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): 這邊認定 esp_sysid 會等於 mavlink sysid + poll time = 0.2305 * grant_bytes + 58.094 (R平方為0.979 高線性相關) + 硬體產品系列 (Product Family) : XBP9B-DM 晶片世代: XBee PRO 900HP 200K 運作頻段: 900 MHz RF。 @@ -295,7 +302,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): MAX_PAYLOAD_PER_FRAME = 80 CHUNK_SEND_INTERVAL_SEC = 0.005 - + class Esp32DeviceInfo: def __init__(self, system_id, address_64, last_hello_time): @@ -304,15 +311,17 @@ 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_poll_time = 0.0 # 目前沒用到 self.last_done_time = 0.0 # 最後送出Done的時間 self.received_len = 0 # 收到封包累計 + self.poll_miss_count = 0 # poll 封包連續丟失的累計數 + def __init__(self, at_handler: "ATCommandHandler" = None): super().__init__(at_handler) - self.max_discovery_window_ms = 220 - self.is_discovery_phase = False + self.max_discovery_window_ms = 0.22 + self.is_discovery_phase = False # 只是 show 狀態 用不太到 self.esp32_address_mapping = {} # bytes[Addr64] : Esp32DeviceInfo self.operator_busy = False @@ -326,21 +335,21 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.event_loop: Optional[asyncio.AbstractEventLoop] = None self.serial_baudrate = 115200 - self.gcs_transmit_queue: deque[bytearray] = deque() + self.udp_transmit_queue: deque[bytearray] = deque() self.poll_scheduler_state = pollStrategy.PollSchedulerState() self.command_pending_event: Optional[asyncio.Event] = None self.poll_done_event: Optional[asyncio.Event] = None self.current_poll_address_64: Optional[bytes] = None - + self.current_poll_missing: bool = False self.current_poll_time = 0 # 本次 poll 丟出的時間紀錄 用來運算傳輸效能 self.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間 self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的 - self.mavPack_interval_timeout = 150 # ms,poll 期間 MAVLink/DONE 最大閒置間隔 - self.discovery_interval_seconds = 200.0 # 每次做 discovery 程序的間隔時間 + self.mavPack_interval_timeout = 150 # (msec) poll 期間 每筆 MAVLink 資料 /DONE 最大閒置間隔 + self.discovery_interval_seconds = 200.0 # (sec) 每次做 discovery 程序的間隔時間 self.device_offline_timeout = self.discovery_interval_seconds * 2 # 遠端沒有回應會被踢出 超時時限 - self.operator_tick_interval_seconds = 0.03 # + self.operator_tick_interval_seconds = 0.01 # self.guard_milliseconds = 50 # POLL DONE 的保底時間間隔 self.pending_manual_discovery = False @@ -348,6 +357,11 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0xFE) + self._pending_tx_count = 0 + self._last_tx_sent_time = 0.0 + self._next_gcs_frame_id = 0x10 + self._rf_busy_timeout_sec = 2 + # ---- 注入與設定 ---- def set_writer(self, writer: Callable[[bytes], None]) -> None: @@ -361,6 +375,13 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.serial_baudrate = baudrate def _ensure_async_primitives(self) -> None: + """ + 因為 FrameProcessor 被歸屬在 SerialHandler 裡面 + SerialHandler 物件在初始化時 是在 main thread + + 然而 以下兩個 Event 物件是在分支 thread 運作的 + 所以這裡不能放到 __init__ 中 + """ if self.command_pending_event is None: self.command_pending_event = asyncio.Event() if self.poll_done_event is None: @@ -371,12 +392,145 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self._ensure_async_primitives() self.command_pending_event.set() + # ---- 對外契約 ---- + + def process_outgoing(self, data: bytes) -> bytes: + """這裡不在用 Base 的方式送出""" + self.enqueue_serial_transmit(data) + self.wake_operator() + return b'' + + # ---- UDP 到 Serial 佇列 ---- + + # 把 UDP 的資料先放到 udp_transmit_queue + def enqueue_serial_transmit(self, payload: bytes) -> None: + if not payload: + return + self.udp_transmit_queue.append(bytearray(payload)) + + # 只是 show 狀態 + def get_gcs_queue_byte_count(self) -> int: + return sum(len(packet) for packet in self.udp_transmit_queue) + + # 只是 show 狀態 + def get_gcs_queue_packet_count(self) -> int: + return len(self.udp_transmit_queue) + + # ---- UDP 佇列透過 RF 模組 TX ---- + + # 判斷RF模組是不是已經完成送出 + @property + def rf_busy(self) -> bool: + if self._pending_tx_count > 0: + if time.time() - self._last_tx_sent_time > self._rf_busy_timeout_sec: + logger.warning("RF busy timeout, force idle") + self._pending_tx_count = 0 + return self._pending_tx_count > 0 + return False + + # 動態分配 frame_id 從 0x10~0xF0 + def _alloc_tx_frame_id(self) -> int: + frame_id = self._next_gcs_frame_id + self._next_gcs_frame_id += 1 + if self._next_gcs_frame_id > 0xF0: + self._next_gcs_frame_id = 0x10 + return frame_id + + # def _notify_tx_sent(self, frame_id: int) -> None: + # if frame_id == 0x00: + # return + # self._pending_tx_count += 1 + # self._last_tx_sent_time = time.time() + + def _handle_tx_status(self, frame: bytes) -> None: + if len(frame) < 9 : + return + frame_id = frame[4] + delivery = frame[8] + + if frame_id == 0x02: # poll pack 的回傳處理 + # 若封包丟失 則設立旗標並停下等待迴圈 + if delivery != 0x00: + self.current_poll_missing = True + self.poll_done_event.set() + + + elif (frame_id >= 0x10) and (frame_id > 0xF0): # UDP 上傳到 RF 的處理 + if self._pending_tx_count > 0: + self._pending_tx_count -= 1 + + if delivery != 0x00: + logger.warning( + f"TX failed fid=0x{frame[4]:02X} delivery=0x{delivery:02X}" + ) + + # 把 udp_transmit_queue 的 mavlink 封包依照大小打包出來 + def _pop_flush_batch(self, max_bytes: int) -> List[bytes]: + if not self.udp_transmit_queue: + return [] + + batch: List[bytes] = [] + total_bytes = 0 + + while self.udp_transmit_queue: + next_packet = bytes(self.udp_transmit_queue[0]) + packet_length = len(next_packet) + + # batch 沒東西 就轉移一筆 udp_transmit_queue 的資料進去 + # 如果有就算一下這筆加進去會不會超出 + if (not batch) or (total_bytes + packet_length <= max_bytes): + batch.append(self.udp_transmit_queue.popleft()) + total_bytes += packet_length + else: + break + + return batch + + # # 把數據封裝好以後 交給 serial_writer 排程送出 + # async def _send_udp_packet(self, packet: bytes) -> None: + # sent_offset = 0 + # while sent_offset < len(packet): + # chunk_end = min( + # sent_offset + self.MAX_PAYLOAD_PER_FRAME, + # len(packet), + # ) + # chunk = packet[sent_offset:chunk_end] + # sent_offset = chunk_end + # frame_id = self._alloc_tx_frame_id() + # self.serial_writer( + # self._encapsulate( + # chunk, + # dest_addr16 = self.DEST_ADDR16_BRAODCAST, + # dest_addr64 = self.DEST_ADDR64_BRAODCAST, + # frame_id=frame_id)) + # self._notify_tx_sent(frame_id) + # await asyncio.sleep(self.CHUNK_SEND_INTERVAL_SEC) + + # 處理從 UDP 來的資訊 丟給 Serial + async def flush_udp_transmit_queue(self) -> None: + if self.serial_writer is None: + logger.warning("GCS flush skipped: serial writer not ready") + return + + for packet in self._pop_flush_batch(self.MAX_BYTES_PER_FLUSH): + # await self._send_udp_packet(packet) + frame_id = self._alloc_tx_frame_id() + self.serial_writer( + self._encapsulate( + packet, + dest_addr16 = self.DEST_ADDR16_BRAODCAST, + dest_addr64 = self.DEST_ADDR64_BRAODCAST, + frame_id=frame_id)) + # self._notify_tx_sent(frame_id) + self._pending_tx_count += 1 + self._last_tx_sent_time = time.time() + # ---- 拆幀分派 ---- def _dispatch_frame(self, frame: bytes) -> Optional[bytes]: frame_type = frame[3] - if frame_type == self.FRAME_TYPE_RX_PACKET: + if frame_type == ATCommandHandler.FRAME_TYPE_RX_PACKET: payload, sender_address_64 = self._decapsulate(frame) if payload.startswith(self.HELLO_HEADER) and len(payload) == 5: @@ -394,18 +548,19 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.last_recieve_mavlink = time.time() return payload - if frame_type == self.FRAME_TYPE_AT_RESPONSE: + if frame_type == ATCommandHandler.FRAME_TYPE_AT_RESPONSE: if self.at_handler is not None: self.at_handler.handle_frame(frame) return None - if frame_type == self.FRAME_TYPE_TX_STATUS: - # length = (frame[1] << 8) | frame[2] - # logger.debug( - # f"TX Status raw={frame.hex()}, api_len={length}, " - # f"fid=0x{frame[4]:02X}, dest16=0x{(frame[5]<<8)|frame[6]:04X}, " - # f"retry={frame[7]}, delivery={frame[8]}, discovery={frame[9]}" - # ) + if frame_type == ATCommandHandler.FRAME_TYPE_TX_STATUS: + self._handle_tx_status(frame) + length = (frame[1] << 8) | frame[2] # for debug + logger.info( + f"TX Status raw={frame.hex()}, api_len={length}, " + f"fid=0x{frame[4]:02X}, dest16=0x{(frame[5]<<8)|frame[6]:04X}, " + f"retry={frame[7]}, delivery={frame[8]}, discovery={frame[9]}" + ) # for debug return None logger.warning(f"Unknown XBee frame type: 0x{frame_type:02X}") @@ -415,7 +570,11 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): def pack_discovery(self) -> bytes: logger.debug(f"pack discovery") - return self._encapsulate(self.DISC_HEADER, dest_addr16 = self.DEST_ADDR16_BRAODCAST, dest_addr64 = self.DEST_ADDR64_BRAODCAST , frame_id=0x00) + return self._encapsulate( + self.DISC_HEADER, + dest_addr16 = self.DEST_ADDR16_BRAODCAST, + dest_addr64 = self.DEST_ADDR64_BRAODCAST, + frame_id=0x05) # 處理每個裝置回傳的 Hello 訊息 def handle_hello_report(self, payload: bytes, sender_address_64: bytes) -> None: @@ -448,24 +607,21 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): # f"address_64={sender_address_64.hex()}" # ) - - # ---- POLL 封裝 / DONE 處理 ---- - def pack_poll(self, target_address_64: bytes, grant_bytes: int = 0) -> Optional[bytes]: - # logger.debug(f"pack poll {target_address_64.hex()} / {grant_bytes}") - remote_device = self.esp32_address_mapping.get(target_address_64) - if remote_device is None: - return None + # ---- POLL 封裝 / DONE 處理 ---- + def pack_poll(self, remote_device, grant_bytes: int = 0) -> Optional[bytes]: + # logger.debug(f"pack poll {remote_device.address_64.hex()} / {grant_bytes}") poll_payload = self.POLL_HEADER + struct.pack( '>BH', remote_device.system_id, grant_bytes ) - remote_device.received_len = 0 + return self._encapsulate( poll_payload, dest_addr64=remote_device.address_64, dest_addr16=self.DEST_ADDR16_UNICAST, frame_id=0x02, + options=0x01 # xbee auto retry ) def handle_done_report(self, payload: bytes, sender_address_64: bytes) -> None: @@ -479,7 +635,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): # 紀錄傳輸狀況 以利後續算速率 transmit_time = int((time.time() - self.current_poll_time)*1000) # ms vehicle = vehicle_registry.get(system_id) - logger.debug(f"Done! payload size:{remote_device.received_len} / time(ms):{transmit_time:.2f}") + + # logger.debug(f"Done! payload size:{remote_device.received_len} / time(ms):{transmit_time:.2f}") # if vehicle and not vehicle.rf_module: # vehicle.rf_module.poll_total_size -= vehicle.rf_module.poll_data_sizes[0] @@ -500,77 +657,6 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): ): self.poll_done_event.set() - # ---- UDP到Serial 佇列 ---- - - def enqueue_gcs_transmit(self, payload: bytes) -> None: - if not payload: - return - self.gcs_transmit_queue.append(bytearray(payload)) - - # 只是 show 狀態 - def get_gcs_queue_byte_count(self) -> int: - return sum(len(packet) for packet in self.gcs_transmit_queue) - - # 只是 show 狀態 - def get_gcs_queue_packet_count(self) -> int: - return len(self.gcs_transmit_queue) - - # 把 gcs_transmit_queue 的 mavlink 封包依照大小打包出來 - def _pop_flush_batch(self, max_bytes: int) -> List[bytes]: - if not self.gcs_transmit_queue: - return [] - - batch: List[bytes] = [] - total_bytes = 0 - - while self.gcs_transmit_queue: - next_packet = bytes(self.gcs_transmit_queue[0]) - packet_length = len(next_packet) - - if not batch: - batch.append(self.gcs_transmit_queue.popleft()) - total_bytes = packet_length - if packet_length > max_bytes: - break - continue - - if total_bytes + packet_length <= max_bytes: - batch.append(self.gcs_transmit_queue.popleft()) - total_bytes += packet_length - else: - break - - return batch - - # 把數據封裝好以後 交給 serial_writer 排程送出 - async def _send_gcs_packet(self, packet: bytes) -> None: - # 重複判斷 - # if self.serial_writer is None: - # return - - sent_offset = 0 - while sent_offset < len(packet): - chunk_end = min( - sent_offset + self.MAX_PAYLOAD_PER_FRAME, - len(packet), - ) - chunk = packet[sent_offset:chunk_end] - sent_offset = chunk_end - self.serial_writer(self._encapsulate(chunk)) - await asyncio.sleep(self.CHUNK_SEND_INTERVAL_SEC) - - # 把上面兩個步驟打包起來 處理從 UDP 來的資訊 丟給 Serial - async def flush_gcs_transmit_queue( - self, - max_bytes: int = MAX_BYTES_PER_FLUSH, - ) -> None: - if self.serial_writer is None: - logger.warning("GCS flush skipped: serial writer not ready") - return - - for packet in self._pop_flush_batch(max_bytes): - await self._send_gcs_packet(packet) - # ---- POLL 排程輔助 ---- # 計算下次要 poll 的對象跟大小 @@ -610,6 +696,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): def _should_run_discovery(self) -> bool: # 條件1. 目前沒有任何遠端ESP裝置被紀錄 或者 手動啟動 if (not self.esp32_address_mapping) or (self.pending_manual_discovery): + # logger.debug(f"mark D : {(time.time() - self.last_discovery_time) > (1 + 5 * self.max_discovery_window_ms)}") return (time.time() - self.last_discovery_time) > (1 + 5 * self.max_discovery_window_ms) # 條件2. 每個固定週期 會做一次 return (time.time() - self.last_discovery_time) > self.discovery_interval_seconds @@ -683,6 +770,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): }) return { "operator_busy": self.operator_busy, + "rf_busy": self.rf_busy, + "pending_tx_count": self._pending_tx_count, "is_discovery_phase": self.is_discovery_phase, "gcs_queue_bytes": self.get_gcs_queue_byte_count(), "gcs_queue_packets": self.get_gcs_queue_packet_count(), @@ -736,7 +825,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): # logger.debug(f"now speed: {speed:.2f} / total_size: {vehicle.rf_module.poll_total_size}") - + # 2. 移除沒反應 dongle self._prune_stale_devices() @@ -746,16 +835,20 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): return # 4. (最優先) 處理從 UDP 過來的封包 並且透過 Serial 送出 - if self.gcs_transmit_queue: - await self.flush_gcs_transmit_queue() + if self.udp_transmit_queue: + await self.flush_udp_transmit_queue() return - # 5. 檢測要不要跑 discovery 程序 + # 5. RF 仍在傳輸,略過 discovery / poll + if self.rf_busy: + return + + # 6. 檢測要不要跑 discovery 程序 if self._should_run_discovery(): await self._run_discovery() return - # 6. POLL 程序 (若有手動先執行 若無自動則策略決策) + # 7. POLL 程序 (若有手動先執行 若無自動則策略決策) if self.pending_manual_poll is not None: target_system_id, grant_bytes = self.pending_manual_poll self.pending_manual_poll = None @@ -816,31 +909,40 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): # 這邊把要求封包最小量放在36 是讓載具端正好可以回傳一個有簽章的 HEARTBEAT 封包的大小 # 算是 "探測封包" 這樣讓系統可以更快速的知道載具端殘餘的資料量 - grant_bytes = max(36, min(int(grant_bytes), 65535)) + # 這個別定義在這邊了 讓 pollStategy 好好做自己的事情 + # grant_bytes = max(36, min(int(grant_bytes), 65535)) - poll_frame = self.pack_poll(target_address_64, grant_bytes) - if poll_frame is None: - return + remote_device = self.esp32_address_mapping.get(target_address_64) + if remote_device is None: + return None + + # 打包 poll pack + poll_frame = self.pack_poll(remote_device, grant_bytes) self._ensure_async_primitives() self.operator_busy = True + remote_device.received_len = 0 self.current_poll_address_64 = target_address_64 self.current_poll_time = time.time() + self.current_poll_missing = False self.poll_done_event.clear() self.serial_writer(poll_frame) self.last_recieve_mavlink = time.time() - try: - 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 - await asyncio.sleep(self.guard_milliseconds / 1000.0) + completed = await self._wait_poll_done_with_idle_timeout() + if not completed: + remote_device.remain_bytes -= remote_device.received_len + logger.warning( + f"POLL timeout address_64={target_address_64.hex()} " + f"grant_bytes={grant_bytes}" + ) + # 判斷是否 poll pack 丟失 並紀錄連續丟失 清空 remain_bytes + if self.current_poll_missing: + remote_device.poll_miss_count += 1 + remote_device.remain_bytes = 0 + self.operator_busy = False + self.current_poll_address_64 = None + # await asyncio.sleep(self.guard_milliseconds / 1000.0) # poll 程序中 只要目標的載具還有在傳資訊 就不會終止資料的接收 async def _wait_poll_done_with_idle_timeout(self) -> bool: @@ -849,6 +951,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): while not self.poll_done_event.is_set(): if time.time() - self.last_recieve_mavlink >= idle_timeout_sec: + # logger.debug(f"mark A poll timeout") return False try: await asyncio.wait_for( @@ -861,7 +964,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): def stop_operator(self) -> None: self.operator_running = False - self.gcs_transmit_queue.clear() + self.udp_transmit_queue.clear() self.wake_operator() @@ -881,6 +984,12 @@ class ATCommandHandler: """ FRAME_TYPE_AT_COMMAND = 0x08 + FRAME_TYPE_AT_RESPONSE = 0x88 + + FRAME_TYPE_TX_REQUEST = 0x10 + FRAME_TYPE_TX_STATUS = 0x8B + + FRAME_TYPE_RX_PACKET = 0x90 def __init__(self, serial_port: str): self.send_command_time = 0.0 # for dev @@ -893,9 +1002,8 @@ class ATCommandHandler: # 可擴展其他 AT 指令 } - # ---- 發送端 ---- + # ---- 發送端 ---- # 這個不應該出現在這 越權了 def set_writer(self, writer): - """由 SerialHandler 注入實際寫入 serial 的 callable""" self.writer = writer def send_command(self, request: ATRequest): @@ -1083,15 +1191,9 @@ class UDPHandler(asyncio.DatagramProtocol): logger.warning("Serial handler not set, dropping UDP packet") return - if self.serial_mode == SerialMode.XBEEAPI_espv1: - processor = self.serial_handler.processor - # if isinstance(processor, XBeeFrameProcessor_ESPv1): # 預設唯一綁定 少一個 if 多一點效率 - processor.enqueue_gcs_transmit(data) - processor.wake_operator() - return - processed_data = self.serial_handler.processor.process_outgoing(data) - self.serial_handler.transport.write(processed_data) + if processed_data: + self.serial_handler.transport.write(processed_data) # ====================== Serial Manager ===================== @@ -1479,8 +1581,12 @@ if __name__ == '__main__': 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) + time.sleep(1) + processor.request_discovery() + time.sleep(1) + 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()) @@ -1494,7 +1600,7 @@ if __name__ == '__main__': # sm.send_at_command(1, rssi_request) # time.sleep(1) - time.sleep(10) + time.sleep(120) sm.remove_serial_link(serial_id) sm.shutdown()