diff --git a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py index 7ccef8d..0b74394 100644 --- a/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py +++ b/src/fc_network_adapter/fc_network_adapter/mavlinkVehicleView.py @@ -37,6 +37,7 @@ class ComponentType(Enum): class RFModuleType(Enum): """RF模組類型""" XBEE = "xbee" + XBEE_ESP = "xbee_esp" UDP = "udp" TCP = "tcp" OTHER = "other" diff --git a/src/fc_network_adapter/fc_network_adapter/serialManager.py b/src/fc_network_adapter/fc_network_adapter/serialManager.py index 53714e8..d6e538f 100644 --- a/src/fc_network_adapter/fc_network_adapter/serialManager.py +++ b/src/fc_network_adapter/fc_network_adapter/serialManager.py @@ -311,7 +311,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self.max_discovery_window_ms = 220 self.is_discovery_phase = False - self.esp32_address_mapping = {} + self.esp32_address_mapping = {} # bytes[Addr64] : Esp32DeviceInfo self.operator_busy = False self.operator_running = False @@ -326,9 +326,10 @@ 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.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間 - self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的 + 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 # 遠端沒有回應會被踢出 超時時限 @@ -427,7 +428,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): if vehicle:=vehicle_registry.get(system_id): if not vehicle.rf_module: - vehicle.rf_module = RFModule(RFModuleType.XBEE) + vehicle.rf_module = RFModule(RFModuleType.XBEE_ESP) def pack_poll(self, target_address_64: bytes, grant_bytes: int = 0) -> Optional[bytes]: remote_device = self.esp32_address_mapping.get(target_address_64) @@ -549,7 +550,6 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): poll_devices = [ pollStrategy.PollDevice( address_64=address_64, - # system_id=device.system_id, remain_bytes=device.remain_bytes, last_done_time=device.last_done_time, ) @@ -753,22 +753,6 @@ 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: @@ -785,6 +769,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): self._ensure_async_primitives() self.operator_busy = True self.current_poll_address_64 = target_address_64 + self.current_poll_time = time.time() self.poll_done_event.clear() self.serial_writer(poll_frame) self.last_recieve_mavlink = time.time() @@ -801,6 +786,23 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base): 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: + 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 + def stop_operator(self) -> None: self.operator_running = False self.gcs_transmit_queue.clear() 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 292c5d0..62634ae 100644 --- a/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py +++ b/src/fc_network_adapter/fc_network_adapter/utils/pollStrategy.py @@ -12,7 +12,6 @@ from typing import List class PollDevice: """從 esp32AddrMapping 抽出的唯讀快照。""" address_64: bytes - # system_id: int remain_bytes: int last_done_time: float diff --git a/src/unitdev02/unitdev02/devnote.txt b/src/unitdev02/unitdev02/devnote.txt index 1f22c69..c3edb2d 100644 --- a/src/unitdev02/unitdev02/devnote.txt +++ b/src/unitdev02/unitdev02/devnote.txt @@ -27,6 +27,7 @@ python -m fc_network_adapter.tests.test_vehicleStatusPublisher python -m fc_network_adapter.tests.test_ringBuffer python -m fc_network_adapter.fc_network_adapter.mainOrchestrator +python -m fc_network_adapter.fc_network_adapter.serialManager python -m someotherpkg.src.example_takeoff_land python -m someotherpkg.src.example_change_mode diff --git a/src/unitdev02/unitdev02/esp32.py b/src/unitdev02/unitdev02/esp32.py new file mode 100644 index 0000000..243b44a --- /dev/null +++ b/src/unitdev02/unitdev02/esp32.py @@ -0,0 +1,641 @@ +from machine import UART +import time +import struct +import gc + +# ========================================================= +# ESP32 / MicroPython:UAV 端 XBee <-> Flight Controller Bridge +# Packet-size TDMA + DONE + SYSID / GCS address 自動學習版 +# +# 本版已刪除「一般 DISC / 強制 DSCF」雙模式。 +# 現在只保留一種 discovery 封包,名稱統一叫 DISC。 +# +# 重要: +# - DISC 的功能等同於原本的 DSCF,也就是「強制 discovery」。 +# - UAV 收到 DISC 後,只要已知 MY_SYSID 與 DEST_64,就會重新回 HELO。 +# - 若收到 DISC 時 MY_SYSID 尚未學到,會設定 DISC_PENDING; +# 之後一旦從飛控 MAVLink 學到 SYSID,就自動補送 HELO。 +# +# 封包格式: +# GCS -> UAV: +# DISC +# POLL + target_sysid(1) + grant_bytes(2) +# +# UAV -> GCS: +# HELO + sysid(1) +# DONE + sysid(1) + sent_len(2) + remain_len(2) +# +# 注意: +# - HELO 不帶 remain。 +# - DONE 保留 remain。 +# - MY_SYSID 不寫死,由飛控 MAVLink 自動學。 +# - DEST_64 不寫死,由 GCS 封包的 XBee 0x90 src64 自動學。 +# ========================================================= + + +# ================= 設定區 ================= +FC_BAUDRATE = 115200 +XB_BAUDRATE = 115200 + +DEST_64 = None +MY_SYSID = None + +SYSID_LEARN_CONFIRM_COUNT = 3 +_sysid_candidate = None +_sysid_candidate_count = 0 + +# 若收到 DISC 時還沒學到 MY_SYSID,就先記住。 +# 等之後從 FC MAVLink 學到 MY_SYSID 後,自動補送 HELO。 +DISC_PENDING = False + +XBEE_MAX_PAYLOAD = 100 + +# Discovery:只保留 DISC;功能等同原本 DSCF,收到後強制回 HELO。 +DISC_MAGIC = b'DISC' + +# UAV -> GCS hello: +# HELO + sysid(1) +HELLO_MAGIC = b'HELO' + +HELLO_MIN_INTERVAL_MS = 1000 +_last_hello_ms = 0 + +# 保留此變數只作狀態記錄;本版 DISC 會 force=True,所以不會被 HELLO_SENT 擋掉。 +HELLO_SENT = False + +POLL_MAGIC = b'POLL' +DEFAULT_GRANT_BYTES = 600 + +DONE_MAGIC = b'DONE' + +MAX_BUF_SIZE = 6144 +READ_FC_BYTES = 250 + +FC_UART_ID = 1 +FC_TX_PIN = 32 +FC_RX_PIN = 33 + +XB_UART_ID = 2 +XB_TX_PIN = 25 +XB_RX_PIN = 26 +# ========================================= + + +uart_fc = UART(FC_UART_ID, baudrate=FC_BAUDRATE, tx=FC_TX_PIN, rx=FC_RX_PIN, rxbuf=4096) +uart_xb = UART(XB_UART_ID, baudrate=XB_BAUDRATE, tx=XB_TX_PIN, rx=XB_RX_PIN, rxbuf=4096) + +tx_buf = bytearray() # FC -> GCS 等待 TDMA poll 的 MAVLink stream +rx_buf = bytearray() # XBee API frame parser buffer + + +# ========================================================= +# 基本工具函式 +# ========================================================= + +def get_checksum(data): + return 0xFF - (sum(data) & 0xFF) + + +def is_valid_addr64(addr): + if addr is None: + return False + if len(addr) != 8: + return False + if addr == b'\x00\x00\x00\x00\x00\x00\x00\x00': + return False + return True + + +def learn_dest64_from_src64(src64): + """ + UAV 收到 GCS 的 XBee 0x90 frame 時, + 0x90 裡面的 src64 就是 GCS / coordinator XBee 的 64-bit address。 + """ + global DEST_64 + global HELLO_SENT + global DISC_PENDING + + if not is_valid_addr64(src64): + return False + + src64 = bytes(src64) + + if DEST_64 is None: + DEST_64 = src64 + HELLO_SENT = False + return True + + if DEST_64 != src64: + DEST_64 = src64 + HELLO_SENT = False + DISC_PENDING = False + return True + + return False + + +def build_api_tx_frame(payload): + """ + 建立 XBee API frame type 0x10,把 payload 送到 DEST_64。 + 若 DEST_64 尚未學到,回傳 None。 + """ + if DEST_64 is None: + return None + + frame_content = bytearray() + frame_content.append(0x10) + frame_content.append(0x00) + frame_content.extend(DEST_64) + frame_content.extend(b'\xFF\xFE') + frame_content.append(0x00) + frame_content.append(0x00) + frame_content.extend(payload) + + length = len(frame_content) + + packet = bytearray() + packet.append(0x7E) + packet.append((length >> 8) & 0xFF) + packet.append(length & 0xFF) + packet.extend(frame_content) + packet.append(get_checksum(frame_content)) + + return packet + + +def send_to_xbee_chunked(payload): + """ + 把 payload 依 XBEE_MAX_PAYLOAD 切成多個 XBee API TX frame 送出。 + """ + if DEST_64 is None: + return False + + total_len = len(payload) + sent_len = 0 + + while sent_len < total_len: + end_len = min(sent_len + XBEE_MAX_PAYLOAD, total_len) + chunk = payload[sent_len:end_len] + sent_len = end_len + + pkt = build_api_tx_frame(chunk) + if pkt is None: + return False + + uart_xb.write(pkt) + time.sleep_ms(1) + + return True + + +# ========================================================= +# HELO / DONE 回報 +# ========================================================= + +def send_hello_report(force=False): + """ + UAV 回報自己存在: + HELO + MY_SYSID + + 本版 DISC 預設以 force=True 呼叫,因此每次收到 DISC 都會重新回 HELO。 + """ + global _last_hello_ms + global HELLO_SENT + + if MY_SYSID is None: + return False + if DEST_64 is None: + return False + + if HELLO_SENT and not force: + return False + + now = time.ticks_ms() + if not force: + if time.ticks_diff(now, _last_hello_ms) < HELLO_MIN_INTERVAL_MS: + return False + + payload = HELLO_MAGIC + struct.pack('>B', MY_SYSID) + pkt = build_api_tx_frame(payload) + if pkt is None: + return False + + uart_xb.write(pkt) + _last_hello_ms = now + HELLO_SENT = True + + time.sleep_ms(1) + return True + + +def try_send_pending_hello(): + """ + 若先收到 DISC、後學到 MY_SYSID,這裡會自動補送 HELO。 + """ + global DISC_PENDING + + if DISC_PENDING and MY_SYSID is not None and DEST_64 is not None: + delay_ms = 20 + (((MY_SYSID * 37) + (time.ticks_ms() & 0xFF)) % 180) + time.sleep_ms(delay_ms) + + if send_hello_report(force=True): + DISC_PENDING = False + + +def send_done_report(sent_len, remain_len): + """ + DONE payload: + b'DONE' + sysid(1) + sent_len(2) + remain_len(2) + """ + if MY_SYSID is None: + return False + if DEST_64 is None: + return False + + if sent_len > 65535: + sent_len = 65535 + if remain_len > 65535: + remain_len = 65535 + + payload = DONE_MAGIC + struct.pack('>BHH', MY_SYSID, sent_len, remain_len) + pkt = build_api_tx_frame(payload) + if pkt is None: + return False + + uart_xb.write(pkt) + time.sleep_ms(1) + return True + + +# ========================================================= +# MAVLink frame 解析與 MY_SYSID 自動學習 +# ========================================================= + +def find_first_mavlink_magic(buf): + pos_fe = buf.find(b'\xFE') + pos_fd = buf.find(b'\xFD') + + if pos_fe == -1: + return pos_fd + if pos_fd == -1: + return pos_fe + return pos_fe if pos_fe < pos_fd else pos_fd + + +def mavlink_frame_length(buf, start_idx): + if start_idx >= len(buf): + return None + + magic = buf[start_idx] + + if magic == 0xFE: + if len(buf) - start_idx < 2: + return None + payload_len = buf[start_idx + 1] + total_len = payload_len + 8 + if len(buf) - start_idx < total_len: + return None + return total_len + + if magic == 0xFD: + if len(buf) - start_idx < 3: + return None + payload_len = buf[start_idx + 1] + incompat_flags = buf[start_idx + 2] + signed = (incompat_flags & 0x01) != 0 + total_len = payload_len + 12 + (13 if signed else 0) + if len(buf) - start_idx < total_len: + return None + return total_len + + return None + + +def get_mavlink_sysid(buf, start_idx): + if start_idx >= len(buf): + return None + + magic = buf[start_idx] + + if magic == 0xFE: + if len(buf) - start_idx >= 6: + return buf[start_idx + 3] + + if magic == 0xFD: + if len(buf) - start_idx >= 10: + return buf[start_idx + 5] + + return None + + +def learn_my_sysid_from_tx_buf(): + """ + 從 FC -> ESP32 的 tx_buf 中找完整 MAVLink frame, + 並自動學習飛控的 MAVLink SYSID。 + + 若更改飛控 SYSID,建議重新上電 ESP32,讓 MY_SYSID 重新學習。 + """ + global MY_SYSID + global _sysid_candidate + global _sysid_candidate_count + global HELLO_SENT + global tx_buf + + if len(tx_buf) == 0: + return + + idx = 0 + checked = 0 + + while idx < len(tx_buf) and checked < 8: + sub = tx_buf[idx:] + rel_start = find_first_mavlink_magic(sub) + if rel_start == -1: + return + + start = idx + rel_start + frame_len = mavlink_frame_length(tx_buf, start) + if frame_len is None: + return + + sysid = get_mavlink_sysid(tx_buf, start) + + if sysid is not None and sysid > 0: + if MY_SYSID is not None: + return + + if _sysid_candidate == sysid: + _sysid_candidate_count += 1 + else: + _sysid_candidate = sysid + _sysid_candidate_count = 1 + + if _sysid_candidate_count >= SYSID_LEARN_CONFIRM_COUNT: + MY_SYSID = _sysid_candidate + HELLO_SENT = False + return + + idx = start + frame_len + checked += 1 + + +# ========================================================= +# tx_buf MAVLink frame 取出與裁切 +# ========================================================= + +def pop_mavlink_frames_by_quota(quota_bytes): + global tx_buf + + if quota_bytes <= 0 or len(tx_buf) == 0: + return b'' + + start = find_first_mavlink_magic(tx_buf) + if start == -1: + tx_buf = bytearray() + return b'' + + if start > 0: + tx_buf = tx_buf[start:] + + out = bytearray() + + while len(tx_buf) > 0: + if tx_buf[0] not in (0xFE, 0xFD): + start = find_first_mavlink_magic(tx_buf) + if start == -1: + tx_buf = bytearray() + break + tx_buf = tx_buf[start:] + + frame_len = mavlink_frame_length(tx_buf, 0) + if frame_len is None: + break + + if len(out) > 0 and (len(out) + frame_len) > quota_bytes: + break + + if len(out) == 0 and frame_len > quota_bytes: + out.extend(tx_buf[:frame_len]) + tx_buf = tx_buf[frame_len:] + break + + out.extend(tx_buf[:frame_len]) + tx_buf = tx_buf[frame_len:] + + if len(out) >= quota_bytes: + break + + return bytes(out) + + +def trim_tx_buffer_if_needed(): + global tx_buf + + if len(tx_buf) <= MAX_BUF_SIZE: + return + + target_size = MAX_BUF_SIZE // 2 + + while len(tx_buf) > target_size: + start = find_first_mavlink_magic(tx_buf) + if start == -1: + tx_buf = bytearray() + return + + if start > 0: + tx_buf = tx_buf[start:] + continue + + frame_len = mavlink_frame_length(tx_buf, 0) + if frame_len is None: + if len(tx_buf) > target_size: + tx_buf = tx_buf[-target_size:] + return + + tx_buf = tx_buf[frame_len:] + + +def flush_tx_buffer(grant_bytes): + global tx_buf + + if MY_SYSID is None: # TODO 這個重複判斷了 應該刪 不要浪費效率 + return + + if DEST_64 is None: + return + + if grant_bytes <= 0: + send_done_report(0, len(tx_buf)) + return + + data_to_send = pop_mavlink_frames_by_quota(grant_bytes) + sent_len = len(data_to_send) + + if sent_len > 0: + send_to_xbee_chunked(data_to_send) + + send_done_report(sent_len, len(tx_buf)) + + +# ========================================================= +# GCS control payload 解析 +# ========================================================= + +def parse_poll_payload(real_data): + if not real_data.startswith(POLL_MAGIC): + return None, None + + if len(real_data) == 5: + return real_data[4], DEFAULT_GRANT_BYTES + + if len(real_data) == 7: + target_sysid = real_data[4] + grant_bytes = (real_data[5] << 8) | real_data[6] + return target_sysid, grant_bytes + + return None, None + + +def is_discovery_payload(real_data): + return real_data.startswith(DISC_MAGIC) + + +# ========================================================= +# XBee API RX parser +# ========================================================= + +def process_xbee_buffer(): + global rx_buf + global DISC_PENDING + + while True: + start_pos = rx_buf.find(b'\x7E') + if start_pos == -1: + rx_buf = bytearray() + return + + if start_pos > 0: + rx_buf = rx_buf[start_pos:] + + if len(rx_buf) < 3: + return + + pkt_len = (rx_buf[1] << 8) | rx_buf[2] + total_len = pkt_len + 4 + + if pkt_len > 300: + rx_buf = rx_buf[1:] + continue + + if len(rx_buf) < total_len: + return + + packet = rx_buf[3:3 + pkt_len] + checksum_recv = rx_buf[3 + pkt_len] + + if get_checksum(packet) == checksum_recv: + if len(packet) > 0 and packet[0] == 0x90: + # XBee Receive Packet 0x90: + # packet = 90 | src64(8) | src16(2) | options(1) | RF data + if len(packet) >= 12: + src64 = bytes(packet[1:9]) + real_data = packet[12:] + + # ------------------------------------------------- + # DISC:本版唯一 discovery 封包。 + # 功能等同原本 DSCF:收到後強制回 HELO。 + # ------------------------------------------------- + if is_discovery_payload(real_data): + learn_dest64_from_src64(src64) + learn_my_sysid_from_tx_buf() # TODO 加一個 if MY_SYSID is None + + if MY_SYSID is not None and DEST_64 is not None: + delay_ms = 20 + ((MY_SYSID * 37) + (time.ticks_ms() & 0xFF)) % 180 # TODO 這邊會有多少秒延遲? 為啥是37? 不能21或者17? 20 ~ 199 210 / 30000 + time.sleep_ms(delay_ms) + send_hello_report(force=True) + DISC_PENDING = False + else: + DISC_PENDING = True + + # ------------------------------------------------- + # 判斷是不是 POLL + # ------------------------------------------------- + else: # TODO 這邊要用 elif 去判斷 POLL + target_sysid, grant_bytes = parse_poll_payload(real_data) + + if target_sysid is not None: + learn_dest64_from_src64(src64) + learn_my_sysid_from_tx_buf() + try_send_pending_hello() # TODO 這邊三行怪怪的 GCS 下 poll 然後會回應 HELO 這個很怪 然後為何要每次重做上面兩行 拖效率 + + if MY_SYSID is not None and target_sysid == MY_SYSID: + flush_tx_buffer(grant_bytes) + + else: + # 一般 GCS -> FC MAVLink 下行資料 + if DEST_64 is None: # TODO 這個概念不對 隨便一個 poll 就把本來有做過 DISC 的權限搶走了 危險!!! + learn_dest64_from_src64(src64) + + if DEST_64 is not None and src64 == DEST_64: # TODO 從 GCS 到 FC 的 mavlink 封包 切開不要跟 POLL 放在同一個判斷 + uart_fc.write(real_data) + + rx_buf = rx_buf[total_len:] + + else: + rx_buf = rx_buf[1:] + + + +# ========================================================= +# 初始清理 +# ========================================================= + +try: + gc.collect() +except Exception: + pass + + +# ========================================================= +# 主迴圈 +# ========================================================= + +while True: + try: + # ------------------------------------------------- + # FC -> ESP32 tx buffer + # ------------------------------------------------- + if uart_fc.any(): + data = uart_fc.read(READ_FC_BYTES) + if data: + tx_buf.extend(data) + + learn_my_sysid_from_tx_buf() + try_send_pending_hello() + trim_tx_buffer_if_needed() + + # ------------------------------------------------- + # XBee -> ESP32 + # ------------------------------------------------- + if uart_xb.any(): + chunk = uart_xb.read() + if chunk: + rx_buf.extend(chunk) + + process_xbee_buffer() + + time.sleep_ms(1) + + except MemoryError: + tx_buf = bytearray() + rx_buf = bytearray() + + try: + gc.collect() + except Exception: + pass + + time.sleep_ms(10) + + except Exception: + time.sleep_ms(2) diff --git a/src/unitdev02/unitdev02/ntrip_dev.py b/src/unitdev02/unitdev02/ntrip_dev.py new file mode 100644 index 0000000..a24232f --- /dev/null +++ b/src/unitdev02/unitdev02/ntrip_dev.py @@ -0,0 +1,209 @@ +# import socket + +# def get_rtk2go_source_table(): +# host = "rtk2go.com" +# port = 2101 +# +# print(f"正在連線至 {host}:{port} 獲取掛載點列表...\n") +# ... + +from __future__ import annotations + +import socket +import base64 +import time + + +# ── NMEA GGA(告訴 Caster「我在哪」,觸發 RTCM 推流)──────────────── + +def nmea_checksum(body: str) -> str: + """body 不含 '$' 與 '*checksum',例如 'GPGGA,123519,...'""" + value = 0 + for ch in body: + value ^= ord(ch) + return f"{value:02X}" + + +def decimal_to_nmea_dm(deg: float, *, is_latitude: bool) -> tuple[str, str]: + """十進位度 → NMEA 的 (d)dmm.mmmm 與半球字元。""" + if is_latitude: + hemi = "N" if deg >= 0 else "S" + deg_width = 2 + else: + hemi = "E" if deg >= 0 else "W" + deg_width = 3 + deg = abs(deg) + d = int(deg) + m = (deg - d) * 60.0 + return f"{d:0{deg_width}d}{m:07.4f}", hemi + + +def build_gga_sentence(lat_deg: float, lon_deg: float, alt_m: float = 100.0) -> bytes: + """ + 組一筆 $GPGGA 句子(含 checksum),回傳 bytes 可直接 sock.sendall。 + + lat_deg / lon_deg:十進位經緯度(北、東為正)。 + 測試時請改成 mount 服務範圍內的近似位置;正式使用應來自 GNSS 真實定位。 + """ + utc = time.gmtime() + t_str = f"{utc.tm_hour:02d}{utc.tm_min:02d}{utc.tm_sec:02d}.00" + + lat_dm, ns = decimal_to_nmea_dm(lat_deg, is_latitude=True) + lon_dm, ew = decimal_to_nmea_dm(lon_deg, is_latitude=False) + + # quality=1(GPS fix)、8 顆星、HDOP=1.0 僅供測試示意 + body = ( + f"GPGGA,{t_str},{lat_dm},{ns},{lon_dm},{ew}," + f"1,08,1.0,{alt_m:.1f},M,0.0,M,," + ) + sentence = f"${body}*{nmea_checksum(body)}\r\n" + return sentence.encode("ascii") + + +def send_gga(sock: socket.socket, lat_deg: float, lon_deg: float, alt_m: float = 100.0) -> str: + """送出 GGA,回傳可讀句子供列印。""" + payload = build_gga_sentence(lat_deg, lon_deg, alt_m) + sock.sendall(payload) + return payload.decode("ascii").strip() + + +# ── NTRIP 連線 + RTCM 接收 ───────────────────────────────────────── + +def get_rtcm_payload( + host, + port, + mountpoint, + user="", + password="", + target_ids=None, + *, + gga_lat=None, + gga_lon=None, + gga_alt_m=100.0, + gga_interval_sec=5.0, + filter_ids=True, +): + """ + 連線 NTRIP caster 並解析 RTCM3。 + + gga_lat / gga_lon:若兩者皆有設定,握手成功後會週期送出 GGA(許多 Caster 如 SNIP 至少需一次)。 + gga_interval_sec:GGA 重送間隔(秒)。 + filter_ids:True 時只印 target_ids 內的 message;False 印全部(方便學習觀察)。 + """ + if target_ids is None: + target_ids = [] + need_gga = gga_lat is not None and gga_lon is not None + + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.settimeout(10) + + try: + sock.connect((host, port)) + + auth = base64.b64encode(f"{user}:{password}".encode()).decode() + request = f"GET /{mountpoint} HTTP/1.0\r\n" + request += "User-Agent: NTRIP PythonClient\r\n" + request += f"Authorization: Basic {auth}\r\n" + request += "Connection: close\r\n\r\n" + sock.sendall(request.encode()) + + response = sock.recv(4096) + if b"ICY 200 OK" not in response: + print(f"連線失敗: {response.decode(errors='ignore')}") + return + + print(f"已成功連接至 {mountpoint}") + if need_gga: + print( + f"GGA 模式:每 {gga_interval_sec}s 回報位置 " + f"({gga_lat:.6f}, {gga_lon:.6f}),觸發 Caster 推流" + ) + else: + print("未設定 GGA(部分 Caster 如 RTK2GO 仍可推流;SNIP/區域站常需 GGA)") + print(f"RTCM 過濾目標: {target_ids if filter_ids else '(關閉,顯示全部)'}\n") + + last_gga_time = 0.0 + buffer = b"" + + while True: + now = time.time() + if need_gga and (now - last_gga_time >= gga_interval_sec): + line = send_gga(sock, gga_lat, gga_lon, gga_alt_m) + print(f"[GGA] {line}") + last_gga_time = now + + try: + chunk = sock.recv(4096) + except socket.timeout: + continue + + if not chunk: + print("Caster 關閉連線") + break + + buffer += chunk + + while len(buffer) >= 3: + if buffer[0] != 0xD3: + buffer = buffer[1:] + continue + + length = ((buffer[1] & 0x03) << 8) | buffer[2] + total_len = length + 6 + + if len(buffer) < total_len: + break + + packet = buffer[:total_len] + buffer = buffer[total_len:] + msg_id = (packet[3] << 4) | (packet[4] >> 4) + + if not filter_ids or str(msg_id) in target_ids: + print(f"Found ID: {msg_id} | Length: {total_len} bytes") + + except KeyboardInterrupt: + print("\n使用者中斷") + except Exception as e: + print(f"發生錯誤: {e}") + finally: + sock.close() + + +if __name__ == "__main__": + TARGET_LIST = ["1007", "1230", "1074"] + + # ── Case 1:RTK2GO(通常可不送 GGA)──────────────────────────── + # USER = "chiyu1468@hotmail.com" + # HOST = "rtk2go.com" + # PORT = 2101 + # MOUNT = "No1bio_02" + # get_rtcm_payload( + # HOST, PORT, MOUNT, + # user=USER, + # target_ids=TARGET_LIST, + # filter_ids=False, + # ) + + # ── Case 2:區域 Caster(需 GGA,座標請改成服務區內近似點)────── + USER = "uavlab6061" + PW = "iamsupersmart" + HOST = "210.241.63.193" + PORT = 81 + MOUNT = "2020_GNSS" + # 十進位度;請依 SNIP / 基準站服務範圍調 + GGA_LAT = 24.155792000 + GGA_LON = 120.630679000 + + get_rtcm_payload( + HOST, + PORT, + MOUNT, + user=USER, + password=PW, + target_ids=TARGET_LIST, + gga_lat=GGA_LAT, + gga_lon=GGA_LON, + gga_alt_m=100.0, + gga_interval_sec=30.0, + filter_ids=False, + ) diff --git a/src/unitdev02/unitdev02/tryByteDict.py b/src/unitdev02/unitdev02/tryByteDict.py new file mode 100644 index 0000000..36cbd79 --- /dev/null +++ b/src/unitdev02/unitdev02/tryByteDict.py @@ -0,0 +1,33 @@ +import timeit + +# 建立 1000 個元素的字典,避免單一 Key 的特殊雜湊值干擾 +d_int = {i: True for i in range(1000)} +d_bytes = {f"k_{i:06d}".encode(): True for i in range(1000)} + +# 正確命中中間的值 +target_int = 500 +target_bytes_same = f"k_{500:06d}".encode() + +# 將變數傳入 timeit 的全域環境 +setup_env = lambda: None +globals_dict = { + 'd_int': d_int, + 'd_bytes': d_bytes, + 'target_int': target_int, + 'target_bytes_same': target_bytes_same +} + +# 1. int 查詢 +t_int = timeit.timeit("d_int.get(target_int)", globals=globals_dict, number=50000000) + +# 2. bytes 查詢 (同物件:直接傳入現成變數) +t_bytes_same = timeit.timeit("d_bytes.get(target_bytes_same)", globals=globals_dict, number=50000000) + +# 3. bytes 查詢 (真・新物件:每次都在查詢時動態切片生成新物件,強迫重新計算 Hash) +globals_dict['dynamic_source'] = f"k_{500:06d}".encode() + b"dummy" +t_bytes_new = timeit.timeit("d_bytes.get(dynamic_source[:8])", globals=globals_dict, number=50000000) + +print(f"修正後測試結果 (50,000,000 次):") +print(f"int 查詢 : {t_int:.4f} 秒 (1.0x)") +print(f"bytes 查詢 (同物件): {t_bytes_same:.4f} 秒 (約 {t_bytes_same/t_int:.2f}x)") +print(f"bytes 查詢 (新物件): {t_bytes_new:.4f} 秒 (約 {t_bytes_new/t_int:.2f}x)") \ No newline at end of file diff --git a/src/unitdev02/unitdev02/udptest8.py b/src/unitdev02/unitdev02/udptest8.py new file mode 100644 index 0000000..db9011e --- /dev/null +++ b/src/unitdev02/unitdev02/udptest8.py @@ -0,0 +1,1044 @@ +import asyncio +import serial_asyncio +import struct +import serial +import time +import threading +import tkinter as tk +from tkinter import ttk +from collections import deque, defaultdict +from pymavlink import mavutil +import matplotlib.pyplot as plt +import matplotlib.animation as animation +from matplotlib.backends.backend_tkagg import FigureCanvasTkAgg + +# ========================================================= +# GCS 端:XBee Serial <-> UDP Bridge +# Packet-size TDMA + DONE + RSSI + Auto Discovery +# +# 本版已刪除「一般 DISC / 強制 DSCF」雙模式。 +# 現在只保留一種 discovery 封包,名稱統一叫 DISC。 +# +# 重要: +# - DISC 的功能等同於原本的 DSCF,也就是強制 discovery。 +# - GCS 開機自動 discovery 與 GUI Scan UAV 都送 DISC。 +# - UAV 收到 DISC 後會重新回 HELO。 +# +# 封包格式: +# GCS -> UAV: +# DISC +# POLL + target_sysid(1) + grant_bytes(2) +# +# UAV -> GCS: +# HELO + sysid(1) +# DONE + sysid(1) + sent_len(2) + remain_len(2) +# ========================================================= + + +# ================= 多組設備設定 ================= +CONFIGS = [ + {"serial_port": "/dev/ttyUSB0", "udp_port": 14551}, + {"serial_port": "COM15", "udp_port": 14590}, + {"serial_port": "/dev/ttyUSB2", "udp_port": 14553}, + {"serial_port": "/dev/ttyUSB3", "udp_port": 14554}, +] + +SERIAL_BAUDRATE = 115200 +UDP_REMOTE_IP = '127.0.0.1' + +UDP_LISTEN_FIXED_PORT = False + +TARGET_ADDR64 = b'\x00\x00\x00\x00\x00\x00\xFF\xFF' + +INITIAL_ACTIVE_SYSIDS = [] +ACTIVE_SYSIDS = list(INITIAL_ACTIVE_SYSIDS) + +XBEE_ADDR64_TO_SYSID = {} +AUTO_LEARN_XBEE_ADDR = True + +# 只保留 DISC;功能等同原本 DSCF。 +DISC_MAGIC = b'DISC' +HELLO_MAGIC = b'HELO' +POLL_MAGIC = b'POLL' +DONE_MAGIC = b'DONE' + +# ================= Discovery 參數 ================= +AUTO_DISCOVERY_ENABLED = True +AUTO_DISCOVERY_DURATION_SEC = 30.0 +DISCOVERY_INTERVAL_SEC = 1.5 + +# GUI 按 Scan UAV 時送出 DISC burst +DISCOVERY_BURST_COUNT = 5 +DISCOVERY_BURST_GAP_SEC = 0.20 + +# ================= TDMA 參數 ================= +current_grant_bytes = 600 +current_guard_ms = 20 +INIT_GRANT_BYTES = 1200 + +LOSS_TIME_WINDOW_SEC = 5.0 + +RSSI_AVG_WINDOW = 5 + +RSSI_QUERY_ON_DATA_FRAME = False +RSSI_DATA_QUERY_MIN_INTERVAL_SEC = 0.5 + + +# ================= 全域狀態 ================= +uav_states = {} +for sid in ACTIVE_SYSIDS: + uav_states[sid] = {"mode": "NORMAL"} + +rssi_history = defaultdict(lambda: deque(maxlen=5000)) +rssi_raw_windows = defaultdict(lambda: deque(maxlen=RSSI_AVG_WINDOW)) +rssi_time_history = defaultdict(lambda: deque(maxlen=5000)) +rssi_latest_stats = defaultdict(lambda: { + 'raw': None, + 'avg': None, + 'count': 0, + 'time': 0.0, + 'src64': None, +}) + +packet_loss_history = defaultdict(lambda: deque(maxlen=1000)) +packet_loss_time_history = defaultdict(lambda: deque(maxlen=1000)) +mavlink_sequence_tracker = defaultdict(dict) +packet_loss_stats = defaultdict(lambda: {'loss_rate': 0.0, 'total_received': 0, 'total_lost': 0}) + +learned_addr64_to_sysid = {} +learned_sysid_to_addr64 = {} + +tdma_done_events = {} +tdma_last_reports = defaultdict(lambda: {'time': 0.0, 'sent_len': 0, 'remain_len': 0}) + +hello_last_reports = defaultdict(lambda: {'time': 0.0, 'src64': None}) + +ASYNC_LOOP = None +SERIAL_PROTOCOLS = [] + + +# ========================================================= +# 工具函式 +# ========================================================= + +def format_addr64(addr): + if not addr: + return 'unknown' + return ''.join(f'{b:02X}' for b in addr) + + +def ensure_active_sysid(sysid): + if sysid is None: + return + if sysid <= 0: + return + + if sysid not in ACTIVE_SYSIDS: + ACTIVE_SYSIDS.append(sysid) + ACTIVE_SYSIDS.sort() + print(f"[DISCOVERY] New SYSID discovered: {sysid}") + + if sysid not in uav_states: + uav_states[sysid] = {"mode": "NORMAL"} + + if sysid not in tdma_done_events: + try: + tdma_done_events[sysid] = asyncio.Event() + except RuntimeError: + pass + + +def infer_sysid_from_addr64(src64): + if not src64: + return None + if src64 in XBEE_ADDR64_TO_SYSID: + return XBEE_ADDR64_TO_SYSID[src64] + if src64 in learned_addr64_to_sysid: + return learned_addr64_to_sysid[src64] + return None + + +def learn_xbee_source(sysid, src64): + if not AUTO_LEARN_XBEE_ADDR: + return + if sysid is None or src64 is None: + return + if sysid <= 0: + return + + src64 = bytes(src64) + + old_sysid = learned_addr64_to_sysid.get(src64) + if old_sysid is not None and old_sysid != sysid: + print(f"[DISCOVERY] src64 {format_addr64(src64)} changed SYSID {old_sysid} -> {sysid}") + learned_sysid_to_addr64.pop(old_sysid, None) + if old_sysid in ACTIVE_SYSIDS: + ACTIVE_SYSIDS.remove(old_sysid) + uav_states.pop(old_sysid, None) + tdma_done_events.pop(old_sysid, None) + + old_src64 = learned_sysid_to_addr64.get(sysid) + if old_src64 is not None and old_src64 != src64: + learned_addr64_to_sysid.pop(old_src64, None) + + learned_addr64_to_sysid[src64] = sysid + learned_sysid_to_addr64[sysid] = src64 + ensure_active_sysid(sysid) + + +def get_poll_dest_addr64(sysid): + if sysid in learned_sysid_to_addr64: + return learned_sysid_to_addr64[sysid] + + for addr, sid in XBEE_ADDR64_TO_SYSID.items(): + if sid == sysid: + return addr + + return TARGET_ADDR64 + + +def record_rssi(sysid, rssi_positive_db, src64=None): + if sysid is None: + return + + ensure_active_sysid(sysid) + + raw_dbm = -int(rssi_positive_db) + win = rssi_raw_windows[sysid] + win.append(raw_dbm) + avg_dbm = sum(win) / len(win) + now = time.time() + + rssi_history[sysid].append(avg_dbm) + rssi_time_history[sysid].append(now) + rssi_latest_stats[sysid] = { + 'raw': raw_dbm, + 'avg': avg_dbm, + 'count': len(win), + 'time': now, + 'src64': src64, + } + + +def calculate_packet_loss(sysid, compid, current_seq): + ensure_active_sysid(sysid) + + tracker = mavlink_sequence_tracker[sysid] + now = time.time() + + if compid not in tracker: + tracker[compid] = { + 'last_seq': current_seq, + 'history': deque() + } + return 0.0 + + comp_tracker = tracker[compid] + last_seq = comp_tracker['last_seq'] + + if current_seq > last_seq: + expected = current_seq - last_seq + elif current_seq < last_seq: + expected = (255 - last_seq) + current_seq + 1 + else: + return packet_loss_history[sysid][-1] if packet_loss_history[sysid] else 0.0 + + lost = max(0, expected - 1) + + comp_tracker['history'].append((now, expected, lost)) + comp_tracker['last_seq'] = current_seq + + total_expected_all = 0 + total_lost_all = 0 + + for c_data in tracker.values(): + if 'history' in c_data: + while c_data['history'] and (now - c_data['history'][0][0]) > LOSS_TIME_WINDOW_SEC: + c_data['history'].popleft() + + total_expected_all += sum(item[1] for item in c_data['history']) + total_lost_all += sum(item[2] for item in c_data['history']) + + overall_loss_rate = (total_lost_all / total_expected_all) * 100.0 if total_expected_all > 0 else 0.0 + + packet_loss_stats[sysid] = { + 'loss_rate': overall_loss_rate, + 'total_received': total_expected_all - total_lost_all, + 'total_lost': total_lost_all + } + + packet_loss_history[sysid].append(overall_loss_rate) + packet_loss_time_history[sysid].append(now) + return overall_loss_rate + + +def build_api_tx_frame(data: bytes, dest_addr64: bytes, frame_id=0x00) -> bytes: + frame = b'\x10' + struct.pack('>B', frame_id) + dest_addr64 + b'\xFF\xFE\x00\x00' + data + return b'\x7E' + struct.pack('>H', len(frame)) + frame + struct.pack('B', 0xFF - (sum(frame) & 0xFF)) + + +def estimate_tdma_timeout(grant_bytes): + grant_bytes = max(0, int(grant_bytes)) + max_payload = 100 + chunks = max(1, (grant_bytes + max_payload - 1) // max_payload) + + uart_time = (grant_bytes + chunks * 18 + 32) * 10.0 / SERIAL_BAUDRATE + timeout = uart_time + chunks * 0.012 + (current_guard_ms / 1000.0) + 0.35 + return max(0.12, timeout) + + +async def grant_one_uav(serial_protocols, sysid, grant_bytes): + ensure_active_sysid(sysid) + + ev = tdma_done_events.get(sysid) + if ev is None: + ev = asyncio.Event() + tdma_done_events[sysid] = ev + ev.clear() + + for sp in serial_protocols: + if hasattr(sp, 'send_poll'): + sp.send_poll(sysid, grant_bytes) + + try: + await asyncio.wait_for(ev.wait(), timeout=estimate_tdma_timeout(grant_bytes)) + except asyncio.TimeoutError: + print(f"[TDMA] SYSID {sysid} timeout; move to next UAV") + + await asyncio.sleep(current_guard_ms / 1000.0) + + +# ========================================================= +# SerialToUDP:XBee Serial API parser +# ========================================================= + +class SerialToUDP(asyncio.Protocol): + def __init__(self, udp_protocol, serial_port): + self.udp_protocol = udp_protocol + self.serial_port = serial_port + self.buffer = bytearray() + self.gcs_tx_queue = bytearray() + self.transport = None + + self.current_poll_sysid = None + self.current_poll_time = 0.0 + + self.at_frame_id = 0x20 + self.pending_db = {} + self.last_data_db_query_time = 0.0 + + def connection_made(self, transport): + self.transport = transport + if hasattr(self.udp_protocol, 'set_serial_transport'): + self.udp_protocol.set_serial_transport(self) + print(f"[{self.serial_port}] Serial connection established.") + + def send_discovery(self): + """ + 發 DISC。 + 本版 DISC 功能等同原本 DSCF,UAV 收到後會強制回 HELO。 + """ + api_frame = build_api_tx_frame(DISC_MAGIC, TARGET_ADDR64, 0x00) + self.transport.write(api_frame) + + def send_poll(self, target_sysid, grant_bytes=None): + if grant_bytes is None: + grant_bytes = current_grant_bytes + grant_bytes = max(0, min(int(grant_bytes), 65535)) + + self.current_poll_sysid = target_sysid + self.current_poll_time = time.time() + + poll_payload = POLL_MAGIC + struct.pack('>BH', target_sysid, grant_bytes) + + dest_addr64 = get_poll_dest_addr64(target_sysid) + api_frame = build_api_tx_frame(poll_payload, dest_addr64, 0x00) + + self.transport.write(api_frame) + + def data_received(self, data): + self.buffer.extend(data) + + while True: + try: + start_idx = self.buffer.index(0x7E) + if start_idx > 0: + del self.buffer[:start_idx] + except ValueError: + self.buffer.clear() + return + + if len(self.buffer) < 3: + return + + length = (self.buffer[1] << 8) | self.buffer[2] + full_length = 3 + length + 1 + + if length > 700: + self.buffer.pop(0) + continue + + if len(self.buffer) < full_length: + return + + frame = bytes(self.buffer[:full_length]) + checksum = 0xFF - (sum(frame[3:-1]) & 0xFF) + + if checksum != frame[-1]: + self.buffer.pop(0) + continue + + del self.buffer[:full_length] + self.handle_api_frame(frame) + + def handle_api_frame(self, frame): + frame_type = frame[3] + + if frame_type == 0x90: + src64 = frame[4:12] + rf_data = frame[15:-1] + self.handle_rx_packet(src64, rf_data) + return + + if frame_type == 0x88: + self.handle_at_response(frame) + return + + def handle_rx_packet(self, src64, rf_data): + sysid_hint = infer_sysid_from_addr64(src64) + + if sysid_hint is None and self.current_poll_sysid is not None: + if time.time() - self.current_poll_time <= 2.0: + sysid_hint = self.current_poll_sysid + + result = self.udp_protocol.process_rf_data(rf_data, src64=src64, sysid_hint=sysid_hint) + + if result is not None: + reason, sysid = result + learn_xbee_source(sysid, src64) + self.send_at_command_db(sysid, src64=src64, reason=reason) + return + + if RSSI_QUERY_ON_DATA_FRAME and sysid_hint is not None: + now = time.time() + if now - self.last_data_db_query_time >= RSSI_DATA_QUERY_MIN_INTERVAL_SEC: + self.send_at_command_db(sysid_hint, src64=src64, reason='DATA') + self.last_data_db_query_time = now + + def next_at_frame_id(self): + self.at_frame_id += 1 + if self.at_frame_id > 0xFE: + self.at_frame_id = 0x20 + return self.at_frame_id + + def send_at_command_db(self, sysid, src64=None, reason=''): + try: + frame_id = self.next_at_frame_id() + frame_type = 0x08 + at_command = b'DB' + parameter = b'' + frame_data = struct.pack('>B', frame_type) + struct.pack('>B', frame_id) + at_command + parameter + checksum = 0xFF - (sum(frame_data) & 0xFF) + api_frame = b'\x7E' + struct.pack('>H', len(frame_data)) + frame_data + struct.pack('B', checksum) + + self.pending_db[frame_id] = { + 'sysid': sysid, + 'src64': src64, + 'time': time.time(), + 'reason': reason, + } + + now = time.time() + stale = [fid for fid, info in self.pending_db.items() if now - info['time'] > 2.0] + for fid in stale: + self.pending_db.pop(fid, None) + + self.transport.write(api_frame) + except Exception as e: + print(f"[{self.serial_port}] send ATDB failed: {e}") + + def handle_at_response(self, frame): + if len(frame) < 9: + return + + frame_id = frame[4] + at_command = frame[5:7] + status = frame[7] + + if at_command != b'DB': + return + if status != 0x00 or len(frame) <= 8: + return + + rssi_value = frame[8] + info = self.pending_db.pop(frame_id, None) + + if info is None: + sysid = self.current_poll_sysid + src64 = learned_sysid_to_addr64.get(sysid) + else: + sysid = info.get('sysid') + src64 = info.get('src64') + + if sysid is not None: + record_rssi(sysid, rssi_value, src64=src64) + + def write_to_serial(self, data): + self.gcs_tx_queue.extend(data) + + def flush_gcs_queue(self): + if not self.gcs_tx_queue: + return + + send_limit = min(len(self.gcs_tx_queue), 150) + data_to_send = self.gcs_tx_queue[:send_limit] + self.gcs_tx_queue = self.gcs_tx_queue[send_limit:] + asyncio.create_task(self._async_send_chunks(data_to_send)) + + async def _async_send_chunks(self, data): + try: + max_payload = 80 + sent_len = 0 + while sent_len < len(data): + end_len = min(sent_len + max_payload, len(data)) + chunk = data[sent_len:end_len] + sent_len = end_len + + api_frame = build_api_tx_frame(chunk, TARGET_ADDR64, 0x00) + self.transport.write(api_frame) + await asyncio.sleep(0.01) + except Exception: + pass + + +# ========================================================= +# UDPHandler:UDP / MAVLink / HELO / DONE +# ========================================================= + +class UDPHandler(asyncio.DatagramProtocol): + def __init__(self, udp_port): + self.udp_port = udp_port + self.serial_transport = None + self.transport = None + self.mav_decoder = mavutil.mavlink.MAVLink(None) + + def connection_made(self, transport): + self.transport = transport + + def set_serial_transport(self, serial_transport): + self.serial_transport = serial_transport + + def datagram_received(self, data, addr): + if self.serial_transport: + self.serial_transport.write_to_serial(data) + + def handle_hello_report(self, rf_data, src64=None): + """ + HELO + sysid(1) + """ + if len(rf_data) < 5 or not rf_data.startswith(HELLO_MAGIC): + return None + + try: + sysid = rf_data[4] + + hello_last_reports[sysid] = { + 'time': time.time(), + 'src64': src64, + } + + learn_xbee_source(sysid, src64) + print(f"[HELO] SYSID {sysid}, src64={format_addr64(src64)}") + return sysid + + except Exception: + return None + + def handle_done_report(self, rf_data, src64=None): + if len(rf_data) < 9 or not rf_data.startswith(DONE_MAGIC): + return None + + try: + sysid, sent_len, remain_len = struct.unpack('>BHH', rf_data[4:9]) + + tdma_last_reports[sysid] = { + 'time': time.time(), + 'sent_len': sent_len, + 'remain_len': remain_len, + } + + learn_xbee_source(sysid, src64) + + ev = tdma_done_events.get(sysid) + if ev is not None: + ev.set() + + return sysid + + except Exception: + return None + + def process_rf_data(self, rf_data, src64=None, sysid_hint=None): + hello_sysid = self.handle_hello_report(rf_data, src64=src64) + if hello_sysid is not None: + return ('HELO', hello_sysid) + + done_sysid = self.handle_done_report(rf_data, src64=src64) + if done_sysid is not None: + return ('DONE', done_sysid) + + try: + for byte in rf_data: + msg = self.mav_decoder.parse_char(bytes([byte])) + if msg: + sysid = msg.get_srcSystem() + compid = msg.get_srcComponent() + seq = msg.get_seq() + if sysid == 0: + continue + + if src64 is not None: + learn_xbee_source(sysid, src64) + + calculate_packet_loss(sysid, compid, seq) + + except Exception: + pass + + if self.transport: + self.transport.sendto(rf_data, (UDP_REMOTE_IP, self.udp_port)) + + return None + + +# ========================================================= +# Bridge setup / scheduler / discovery +# ========================================================= + +async def setup_bridge(config): + port, udp = config['serial_port'], config['udp_port'] + + try: + ser = serial.Serial(port, SERIAL_BAUDRATE) + ser.close() + except Exception: + print(f"[{port}] Serial open failed, skip.") + return None + + loop = asyncio.get_running_loop() + udp_handler = UDPHandler(udp) + + local_port = udp if UDP_LISTEN_FIXED_PORT else 0 + await loop.create_datagram_endpoint(lambda: udp_handler, local_addr=('0.0.0.0', local_port)) + + serial_proto = SerialToUDP(udp_handler, port) + await serial_asyncio.create_serial_connection(loop, lambda: serial_proto, port, baudrate=SERIAL_BAUDRATE) + return serial_proto + + +async def auto_discovery_task(serial_protocols): + """ + 程式啟動前 AUTO_DISCOVERY_DURATION_SEC 秒自動送 DISC。 + 本版 DISC 是強制 discovery,UAV 收到後會重新回 HELO。 + """ + if not AUTO_DISCOVERY_ENABLED: + return + + print(f"[DISCOVERY] Auto DISC started for {AUTO_DISCOVERY_DURATION_SEC:.1f} seconds") + + start_t = time.time() + while time.time() - start_t < AUTO_DISCOVERY_DURATION_SEC: + for sp in serial_protocols: + if hasattr(sp, 'send_discovery'): + sp.send_discovery() + await asyncio.sleep(DISCOVERY_INTERVAL_SEC) + + print("[DISCOVERY] Auto DISC stopped. Use GUI Scan UAV for new aircraft.") + + +async def discovery_burst(serial_protocols): + """ + GUI Scan UAV:送 DISC burst。 + """ + print("[DISCOVERY] Manual DISC scan started") + for _ in range(DISCOVERY_BURST_COUNT): + for sp in serial_protocols: + if hasattr(sp, 'send_discovery'): + sp.send_discovery() + await asyncio.sleep(DISCOVERY_BURST_GAP_SEC) + print("[DISCOVERY] Manual DISC scan finished") + + +async def tdma_scheduler(serial_protocols): + print(f"Packet-size TDMA Scheduler Started... 找到 {len(serial_protocols)} 個可用端口") + + for sysid in list(ACTIVE_SYSIDS): + ensure_active_sysid(sysid) + tdma_done_events[sysid] = asyncio.Event() + + while True: + for sp in serial_protocols: + if hasattr(sp, 'flush_gcs_queue'): + sp.flush_gcs_queue() + + await asyncio.sleep(current_guard_ms / 1000.0) + + if not ACTIVE_SYSIDS: + await asyncio.sleep(0.2) + continue + + for sysid in list(ACTIVE_SYSIDS): + ensure_active_sysid(sysid) + + state = uav_states.get(sysid, {"mode": "NORMAL"}) + + if state.get('mode') == 'INITIALIZING': + for _ in range(4): + await grant_one_uav(serial_protocols, sysid, INIT_GRANT_BYTES) + else: + await grant_one_uav(serial_protocols, sysid, current_grant_bytes) + + +async def async_main(): + global ASYNC_LOOP + global SERIAL_PROTOCOLS + + ASYNC_LOOP = asyncio.get_running_loop() + + protocols = await asyncio.gather(*(setup_bridge(cfg) for cfg in CONFIGS)) + valid_protocols = [p for p in protocols if p is not None] + SERIAL_PROTOCOLS = valid_protocols + + if valid_protocols: + asyncio.create_task(auto_discovery_task(valid_protocols)) + asyncio.create_task(tdma_scheduler(valid_protocols)) + else: + print('No valid serial ports found.') + + await asyncio.Future() + + +# ========================================================= +# GUI +# ========================================================= + +def start_gui(): + root = tk.Tk() + root.title('UAV Packet-size TDMA Control Station - Auto Discovery') + root.geometry('1350x900') + + control_frame = tk.Frame(root, width=390, bg='#f0f0f0', padx=20, pady=20) + control_frame.pack(side=tk.LEFT, fill=tk.Y) + + tk.Label(control_frame, text='Packet-size TDMA 控制', font=('Arial', 14, 'bold'), bg='#f0f0f0').pack(pady=10) + + def on_scan_uav(): + if ASYNC_LOOP is not None and SERIAL_PROTOCOLS: + asyncio.run_coroutine_threadsafe(discovery_burst(SERIAL_PROTOCOLS), ASYNC_LOOP) + scan_label.config(text='已送出 DISC 掃描') + else: + scan_label.config(text='尚未建立 serial / asyncio loop') + + scan_btn = tk.Button(control_frame, text='Scan UAV / 發送 DISC', + font=('Arial', 12, 'bold'), bg='#2196F3', fg='white', + command=on_scan_uav) + scan_btn.pack(pady=(5, 5), fill=tk.X) + + scan_label = tk.Label(control_frame, text=f'開機前 {AUTO_DISCOVERY_DURATION_SEC:.0f}s 自動 DISC,之後可手動掃描', + bg='#f0f0f0', font=('Arial', 9), wraplength=330, justify=tk.LEFT) + scan_label.pack(pady=(0, 15), anchor=tk.W) + + grant_var = tk.IntVar(value=current_grant_bytes) + + def on_grant_change(val): + global current_grant_bytes + current_grant_bytes = int(float(val)) + grant_label.config(text=f'每台授權: {current_grant_bytes} bytes') + + grant_slider = ttk.Scale(control_frame, from_=100, to_=1800, orient='horizontal', + variable=grant_var, command=on_grant_change) + grant_slider.pack(fill=tk.X, pady=10) + grant_label = tk.Label(control_frame, text=f'每台授權: {current_grant_bytes} bytes', bg='#f0f0f0') + grant_label.pack() + + guard_var = tk.IntVar(value=current_guard_ms) + + def on_guard_change(val): + global current_guard_ms + current_guard_ms = int(float(val)) + guard_label.config(text=f'切換保護時間: {current_guard_ms} ms') + + tk.Label(control_frame, text='Guard Time', font=('Arial', 11, 'bold'), bg='#f0f0f0').pack(pady=(20, 0)) + guard_slider = ttk.Scale(control_frame, from_=0, to_=100, orient='horizontal', + variable=guard_var, command=on_guard_change) + guard_slider.pack(fill=tk.X, pady=10) + guard_label = tk.Label(control_frame, text=f'切換保護時間: {current_guard_ms} ms', bg='#f0f0f0') + guard_label.pack() + + tk.Label(control_frame, text='群機初始化控制 (INIT)', font=('Arial', 14, 'bold'), bg='#f0f0f0').pack(pady=(30, 10)) + + init_var = tk.IntVar(value=0) + radios = {} + status_labels = {} + + radio_container = tk.Frame(control_frame, bg='#f0f0f0') + radio_container.pack(fill=tk.X) + + def on_radio_change(): + selected = init_var.get() + for sysid in list(ACTIVE_SYSIDS): + ensure_active_sysid(sysid) + if sysid == selected: + uav_states[sysid]['mode'] = 'INITIALIZING' + else: + uav_states[sysid]['mode'] = 'NORMAL' + + tk.Radiobutton(radio_container, text='全體 NORMAL', variable=init_var, value=0, + command=on_radio_change, bg='#f0f0f0', + font=('Arial', 11, 'bold'), fg='blue').pack(anchor=tk.W, pady=(0, 10)) + + def ensure_radio_row(sysid): + if sysid in radios: + return + + frame = tk.Frame(radio_container, bg='#f0f0f0') + frame.pack(fill=tk.X, pady=3) + + rb = tk.Radiobutton(frame, text=f'SYSID {sysid} 專屬載入', + variable=init_var, value=sysid, + command=on_radio_change, + bg='#f0f0f0', font=('Arial', 11)) + rb.pack(side=tk.LEFT) + radios[sysid] = rb + + lbl = tk.Label(frame, text='NORMAL', font=('Arial', 10, 'bold'), fg='green', bg='#f0f0f0') + lbl.pack(side=tk.RIGHT, padx=5) + status_labels[sysid] = lbl + + for sid in ACTIVE_SYSIDS: + ensure_radio_row(sid) + + is_locked = False + + def toggle_tdma_mode(): + nonlocal is_locked + is_locked = not is_locked + + if is_locked: + init_var.set(0) + on_radio_change() + for rb in radios.values(): + rb.config(state=tk.DISABLED) + lock_btn.config(text='點擊解鎖 (允許重新下載參數)', bg='orange') + print('已鎖定進入純 Packet-size TDMA 模式') + else: + for rb in radios.values(): + rb.config(state=tk.NORMAL) + lock_btn.config(text='參數載入完畢,鎖定進入 TDMA', bg='#4CAF50') + print('已解除鎖定,可重新分配特權') + + lock_btn = tk.Button(control_frame, text='參數載入完畢,鎖定進入 TDMA', + font=('Arial', 12, 'bold'), bg='#4CAF50', fg='white', + command=toggle_tdma_mode) + lock_btn.pack(pady=20, fill=tk.X) + + report_title = tk.Label(control_frame, text='最近 HELO / DONE 回報', font=('Arial', 12, 'bold'), bg='#f0f0f0') + report_title.pack(pady=(10, 5)) + + report_container = tk.Frame(control_frame, bg='#f0f0f0') + report_container.pack(fill=tk.X) + + report_labels = {} + hello_labels = {} + + def ensure_report_row(sysid): + if sysid not in hello_labels: + hl = tk.Label(report_container, text=f'SYSID {sysid}: HELO no report', + font=('Arial', 9), bg='#f0f0f0', justify=tk.LEFT) + hl.pack(anchor=tk.W) + hello_labels[sysid] = hl + + if sysid not in report_labels: + dl = tk.Label(report_container, text=f'SYSID {sysid}: DONE sent=0, remain=0', + font=('Arial', 9), bg='#f0f0f0', justify=tk.LEFT) + dl.pack(anchor=tk.W) + report_labels[sysid] = dl + + for sid in ACTIVE_SYSIDS: + ensure_report_row(sid) + + rssi_title = tk.Label(control_frame, text=f'RSSI 最近 {RSSI_AVG_WINDOW} 次平均', font=('Arial', 12, 'bold'), bg='#f0f0f0') + rssi_title.pack(pady=(18, 5)) + + rssi_container = tk.Frame(control_frame, bg='#f0f0f0') + rssi_container.pack(fill=tk.X) + + rssi_labels = {} + + def ensure_rssi_row(sysid): + if sysid in rssi_labels: + return + lbl = tk.Label(rssi_container, text=f'SYSID {sysid}: RSSI avg -- dBm', + font=('Arial', 9), bg='#f0f0f0', justify=tk.LEFT) + lbl.pack(anchor=tk.W) + rssi_labels[sysid] = lbl + + for sid in ACTIVE_SYSIDS: + ensure_rssi_row(sid) + + def update_status_gui(): + now = time.time() + + for sysid in list(ACTIVE_SYSIDS): + ensure_radio_row(sysid) + ensure_report_row(sysid) + ensure_rssi_row(sysid) + + for sysid, lbl in status_labels.items(): + mode = uav_states.get(sysid, {"mode": "NORMAL"}).get("mode", "NORMAL") + if mode == 'INITIALIZING': + lbl.config(text='INIT', fg='orange') + else: + lbl.config(text='NORMAL', fg='green') + + for sysid, lbl in hello_labels.items(): + rep = hello_last_reports[sysid] + age = now - rep['time'] if rep['time'] else 999.0 + age_text = f'{age:.1f}s ago' if age < 99 else 'no report' + addr_text = format_addr64(rep.get('src64')) + lbl.config(text=f"SYSID {sysid}: HELO ({age_text}) src64={addr_text}") + + for sysid, lbl in report_labels.items(): + rep = tdma_last_reports[sysid] + age = now - rep['time'] if rep['time'] else 999.0 + age_text = f'{age:.1f}s ago' if age < 99 else 'no report' + lbl.config(text=f"SYSID {sysid}: DONE sent={rep['sent_len']}, remain={rep['remain_len']} ({age_text})") + + for sysid, lbl in rssi_labels.items(): + st = rssi_latest_stats[sysid] + if st['avg'] is None: + lbl.config(text=f'SYSID {sysid}: RSSI avg -- dBm') + else: + age = now - st['time'] + addr_text = format_addr64(st['src64']) + lbl.config(text=( + f"SYSID {sysid}: avg={st['avg']:.1f} dBm, raw={st['raw']} dBm, " + f"n={st['count']} ({age:.1f}s)\n src64={addr_text}" + )) + + root.after(500, update_status_gui) + + update_status_gui() + + plot_frame = tk.Frame(root) + plot_frame.pack(side=tk.RIGHT, fill=tk.BOTH, expand=True) + + fig, (ax1, ax2) = plt.subplots(2, 1, figsize=(9.5, 8.5), dpi=100) + canvas = FigureCanvasTkAgg(fig, master=plot_frame) + canvas.get_tk_widget().pack(fill=tk.BOTH, expand=True) + + def update_plot(frame): + ax1.clear() + ax2.clear() + + ax1.set_title(f'RSSI avg{RSSI_AVG_WINDOW} from XBee ATDB after HELO / DONE', fontsize=12) + ax1.set_xlim(10, 0) + ax1.set_ylim(-100, -10) + ax1.grid(True, alpha=0.3) + + ax2.set_title('Packet Loss Rate (5s Window)', fontsize=12) + ax2.set_xlim(10, 0) + ax2.set_ylim(0, 100) + ax2.grid(True, alpha=0.3) + + now = time.time() + colors = ['blue', 'red', 'green', 'orange', 'purple', 'brown', 'cyan', 'magenta'] + + try: + sysids = sorted(list(set(list(rssi_history.keys()) + list(packet_loss_history.keys()) + ACTIVE_SYSIDS))) + except RuntimeError: + return + + loss_labels = [] + + for i, sysid in enumerate(sysids): + color = colors[i % len(colors)] + try: + t_hist = list(rssi_time_history[sysid]) + r_hist = list(rssi_history[sysid]) + lt_hist = list(packet_loss_time_history.get(sysid, [])) + l_hist = list(packet_loss_history.get(sysid, [])) + except RuntimeError: + continue + + rssi_recent = [idx for idx, ts in enumerate(t_hist) if now - ts <= 10] + if rssi_recent: + ax1.plot([now - t_hist[idx] for idx in rssi_recent], + [r_hist[idx] for idx in rssi_recent], + label=f'SYSID:{sysid}', color=color, marker='o', markersize=3) + + loss_recent = [idx for idx, ts in enumerate(lt_hist) if now - ts <= 10] + if loss_recent: + loss_t = [now - lt_hist[idx] for idx in loss_recent] + loss_r = [l_hist[idx] for idx in loss_recent] + ax2.plot(loss_t, loss_r, label=f'SYSID:{sysid}', color=color, marker='o', markersize=3) + + if loss_r: + loss_labels.append({ + 'sysid': sysid, + 'y_real': loss_r[-1], + 'x_real': loss_t[-1], + 'color': color + }) + + if loss_labels: + loss_labels = sorted(loss_labels, key=lambda k: k['y_real']) + min_gap = 12.0 + y_positions = [lbl['y_real'] for lbl in loss_labels] + + for j in range(1, len(y_positions)): + if y_positions[j] - y_positions[j - 1] < min_gap: + y_positions[j] = y_positions[j - 1] + min_gap + + if y_positions[-1] > 90: + shift = y_positions[-1] - 90 + y_positions = [y - shift for y in y_positions] + + for j, lbl in enumerate(loss_labels): + sysid = lbl['sysid'] + color = lbl['color'] + real_y = lbl['y_real'] + text_y = y_positions[j] + + ax2.text(0.5, text_y, f'ID:{sysid} ({real_y:.1f}%)', + bbox=dict(boxstyle='round,pad=0.3', facecolor=color, alpha=0.8), + fontsize=10, fontweight='bold', color='white', + horizontalalignment='right', verticalalignment='center') + + if abs(real_y - text_y) > 1.0: + ax2.plot([lbl['x_real'], 0.5], [real_y, text_y], color=color, linestyle=':', alpha=0.6) + + if ax1.lines: + ax1.legend(loc='upper left') + if ax2.lines: + ax2.legend(loc='upper left') + + canvas.draw_idle() + + ani = animation.FuncAnimation(fig, update_plot, interval=1000) + + def on_closing(): + root.quit() + root.destroy() + + root.protocol('WM_DELETE_WINDOW', on_closing) + root.mainloop() + + +# ========================================================= +# Entry point +# ========================================================= + +if __name__ == '__main__': + threading.Thread(target=lambda: asyncio.run(async_main()), daemon=True).start() + start_gui()