diff --git a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py index 0b74394..6184974 100644 --- a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py +++ b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py @@ -182,6 +182,7 @@ class RFStatus: at_response: Optional[str] = None # AT 命令回應 link_quality: Optional[int] = None # 連接品質 timestamp: Optional[float] = None # 時間戳記 + avg_bandwidth: Optional[float] = None # 量測到的傳輸速度 bytes/ms custom_status: Dict[str, Any] = field(default_factory=dict) # 其他自定義狀態 @@ -275,6 +276,11 @@ class RFModule: self.type = rf_type self.status = RFStatus() self.socket_info = SocketInfo() + + def update_avg_bandwidth(self, bandwidth: float, timestamp: Optional[float] = None) -> None: + self.status.avg_bandwidth = bandwidth + if timestamp: + self.status.timestamp = timestamp def update_rssi(self, rssi: int, timestamp: Optional[float] = None) -> None: """更新RSSI""" diff --git a/src/fc_network_adapter/fc_network_adapter/serialManager.py b/src/fc_network_adapter/fc_network_adapter/serialManager.py index afdcb70..b874142 100644 --- a/src/fc_network_adapter/fc_network_adapter/serialManager.py +++ b/src/fc_network_adapter/fc_network_adapter/serialManager.py @@ -218,12 +218,12 @@ class XBeeFrameProcessor_Base(FrameProcessor): return None if frame_type == self.FRAME_TYPE_TX_STATUS: - 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 + # 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}") @@ -293,7 +293,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): MAX_BYTES_PER_FLUSH = 150 MAX_PAYLOAD_PER_FRAME = 80 - CHUNK_SEND_INTERVAL_SEC = 0.01 + CHUNK_SEND_INTERVAL_SEC = 0.005 @@ -308,8 +308,6 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.last_done_time = 0.0 # 最後送出Done的時間 self.received_len = 0 # 收到封包累計 - - def __init__(self, at_handler: "ATCommandHandler" = None): super().__init__(at_handler) @@ -334,7 +332,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): 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_time = 0 # 本次 poll 丟出的時間紀錄 + + self.current_poll_time = 0 # 本次 poll 丟出的時間紀錄 用來運算傳輸效能 self.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間 self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的 @@ -346,6 +345,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.pending_manual_discovery = False self.pending_manual_poll: Optional[tuple[int, Optional[int]]] = None + + self.rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0xFE) # ---- 注入與設定 ---- @@ -410,10 +411,10 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): logger.warning(f"Unknown XBee frame type: 0x{frame_type:02X}") return None - # ---- DISC / POLL 封裝 ---- + # ---- DISC 封裝 / HELLO 處理 ---- def pack_discovery(self) -> bytes: - # logger.debug(f"pack discovery") + 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) # 處理每個裝置回傳的 Hello 訊息 @@ -421,25 +422,37 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): system_id = payload[4] remote_device = self.esp32_address_mapping.get(sender_address_64) + # 若無則新建 if remote_device is None: self.esp32_address_mapping[sender_address_64] = self.Esp32DeviceInfo( system_id, sender_address_64, time.time() ) + + vehicle = vehicle_registry.get(system_id) + if vehicle and not vehicle.rf_module: + 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 + logger.debug( 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: - logger.warning( - f"SYSID duplicated system_id={system_id}, " - f"address_64={sender_address_64.hex()}" - ) + remote_device.last_hello_time = time.time() + + # TODO 檢驗重複 sysid + # elif remote_device.system_id != system_id : + # logger.warning( + # 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_ESP) + # ---- 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 @@ -460,19 +473,23 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): if remote_device is None: return - system_id, sent_length, remain_length = struct.unpack('>BHH', payload[4:9]) remote_device.remain_bytes = remain_length - # 這段是有問題的 因為會有整數封包切割問題 以及載具端的 buffer 存量不足 故回傳的資訊量會與要求的不一致 - # 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}" - # ) + # 紀錄傳輸狀況 以利後續算速率 + 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}") + + # if vehicle and not vehicle.rf_module: + # vehicle.rf_module.poll_total_size -= vehicle.rf_module.poll_data_sizes[0] + # vehicle.rf_module.poll_total_size += remote_device.received_len + # vehicle.rf_module.poll_total_time -= vehicle.rf_module.poll_data_times[0] + # vehicle.rf_module.poll_total_time += transmit_time + # vehicle.rf_module.poll_data_sizes.append(remote_device.received_len) + # vehicle.rf_module.poll_data_times.append(transmit_time) + # logger.debug(f"{vehicle.rf_module.poll_total_size}") - # TODO 傳送速率 - # TODO 累積速率預測 remote_device.received_len = 0 remote_device.last_done_time = time.time() @@ -593,7 +610,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): - return (time.time() - self.last_discovery_time) > 2 + 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 @@ -708,25 +725,37 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): for task in pending: task.cancel() + # 主要迴圈 Main Loop async def _operator_tick(self) -> None: - # 1. 移除沒反應 dongle + # 1. 更新資訊 + # 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) + + # logger.debug(f"now speed: {speed:.2f} / total_size: {vehicle.rf_module.poll_total_size}") + + + + # 2. 移除沒反應 dongle self._prune_stale_devices() - # 2. 忙碌時 略過這次循環 + # 3. 忙碌時 略過這次循環 if self.operator_busy: return - # 3. (最優先) 處理從 UDP 過來的封包 並且透過 Serial 送出 + # 4. (最優先) 處理從 UDP 過來的封包 並且透過 Serial 送出 if self.gcs_transmit_queue: await self.flush_gcs_transmit_queue() return - # 4. 檢測要不要跑 discovery 程序 + # 5. 檢測要不要跑 discovery 程序 if self._should_run_discovery(): await self._run_discovery() return - # 5. POLL 程序 (若有手動先執行 若無自動則策略決策) + # 6. POLL 程序 (若有手動先執行 若無自動則策略決策) if self.pending_manual_poll is not None: target_system_id, grant_bytes = self.pending_manual_poll self.pending_manual_poll = None @@ -738,6 +767,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): return await self._run_one_poll(target_address, grant_bytes if grant_bytes is not None else 0) + # self.at_handler.send_command(self.rssi_request) # 送出 RSSI 請求封包 + def _find_address_by_system_id(self, target_system_id: int) -> Optional[bytes]: for address_64, device in self.esp32_address_mapping.items(): if device.system_id == target_system_id: @@ -761,22 +792,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.mavPack_interval_timeout / 1000.0 - poll_tick = min(0.02, idle_timeout_sec) - - 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, - ) + # async def _wait_poll_done_with_idle_timeout(self) -> bool: + # idle_timeout_sec = self.mavPack_interval_timeout / 1000.0 + # poll_tick = min(0.02, idle_timeout_sec) + + # 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 + # # except asyncio.TimeoutError: + # # continue + # return True # poll 程序 async def _run_one_poll(self, target_address_64: bytes, grant_bytes: int) -> None: @@ -813,7 +844,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): # poll 程序中 只要目標的載具還有在傳資訊 就不會終止資料的接收 async def _wait_poll_done_with_idle_timeout(self) -> bool: - idle_timeout_sec = self.MAX_mavPack_interval_timeout / 1000.0 + idle_timeout_sec = self.mavPack_interval_timeout / 1000.0 poll_tick = min(0.02, idle_timeout_sec / 2) while not self.poll_done_event.is_set(): @@ -852,6 +883,7 @@ class ATCommandHandler: FRAME_TYPE_AT_COMMAND = 0x08 def __init__(self, serial_port: str): + self.send_command_time = 0.0 # for dev self.serial_port = serial_port self.writer = None self.handlers = { @@ -875,6 +907,8 @@ class ATCommandHandler: return self.writer(self._build_at_request(request)) + + self.send_command_time = time.time() # logger.debug( # f"[{self.serial_port}] send AT {request.command.decode(errors='replace')} " # f"(frame_id=0x{request.frame_id:02X})" @@ -902,10 +936,12 @@ class ATCommandHandler: if parsed_at_ack is None: return - if not rx_module_ack.put(parsed_at_ack): - logger.warning( - f"[{self.serial_port}] rx_module_ack overflow, drop {parsed_at_ack.command!r}" - ) + # TODO 本來是打算把 serialManager 的資訊傳給 mavlinkObject 去處理 + # 但是可以考慮 直接寫到 mavlinkVehicleView 裡面 + # if not rx_module_ack.put(parsed_at_ack): + # logger.warning( + # f"[{self.serial_port}] rx_module_ack overflow, drop {parsed_at_ack.command!r}" + # ) self._dispatch(parsed_at_ack) @@ -948,11 +984,9 @@ class ATCommandHandler: def _handle_rssi(self, response: ATResponse): """處理 DB (RSSI) 回應:單 byte 無號值,單位 dBm""" - # print(f"[{self.serial_port}] RSSI = -{data[0]} dBm") # dev - logger.debug(f"[{self.serial_port}] RSSI = -{response.data[0]} dBm") # dev - - - + ack_time = (time.time() - self.send_command_time)*1000 # for dev + # logger.debug(f"{ack_time:.2f}") # for dev + logger.debug(f"[{self.serial_port}] RSSI = -{response.data[0]} dBm / AT ack time:{ack_time}") # dev def _handle_serial_high(self, response: ATResponse): """處理 SH (Serial Number High)""" @@ -1333,12 +1367,14 @@ class serial_manager: logger.error(f"Serial object {serial_id} not found") return False + # TODO 這邊的防呆 應該可以不用 if 有空再改 serial_obj = self.serial_objects[serial_id] - if serial_obj.serial_mode != SerialMode.XBEEAPI2AT: + processor = serial_obj.serial_handler.processor + if not isinstance(processor, XBeeFrameProcessor_Base): logger.error( f"Serial {serial_id} mode is {serial_obj.serial_mode.name}, " - f"AT command only supported in XBEEAPI2AT mode" + f"AT Command NOT supported yet!" ) return False @@ -1432,6 +1468,11 @@ if __name__ == '__main__': UDP_REMOTE_PORT = 14561 sm.create_serial_link(SERIAL_PORT, SERIAL_BAUDRATE, UDP_REMOTE_PORT, SerialMode.XBEEAPI_espv1) + # test_sysid = 20 + # vehicle = vehicle_registry.register(test_sysid) + # if vehicle.rf_module is None: + # vehicle.rf_module = RFModule(RFModuleType.XBEE_ESP) + time.sleep(2) # 等 serial 連線與 operator 啟動 serial_id = 1 @@ -1443,7 +1484,17 @@ if __name__ == '__main__': print(processor.get_status_snapshot()) print(processor.get_gcs_queue_byte_count()) - time.sleep(120) + # rssi_request = ATRequest(command=b'DB', parameter=b'', frame_id=0xFE) + # for i in range(3): + # processor = sm.get_espv1_processor(serial_id) + # if processor is not None: + # deadline = time.time() + 5.0 # 最多等 5 秒,避免永遠卡住 + # while processor.operator_busy and time.time() < deadline: + # time.sleep(0.01) + # sm.send_at_command(1, rssi_request) + # time.sleep(1) + + time.sleep(10) sm.remove_serial_link(serial_id) sm.shutdown()