TempCommit

chiyu
Chiyu Chen 3 weeks ago
parent 2542663e9e
commit 1e07d1fa54

@ -135,24 +135,25 @@ class XBeeFrameProcessor_Base(FrameProcessor):
韌體版本 (Firmware Version): 2014 -> 2 (XBee 3 平台) + 0 (802.15.4 協議代碼) + 14 (次要版本號) 韌體版本 (Firmware Version): 2014 -> 2 (XBee 3 平台) + 0 (802.15.4 協議代碼) + 14 (次要版本號)
""" """
# XBee API frame type # # XBee API frame type
FRAME_TYPE_TX_REQUEST = 0x10 # FRAME_TYPE_TX_REQUEST = 0x10
FRAME_TYPE_AT_RESPONSE = 0x88 # FRAME_TYPE_AT_RESPONSE = 0x88
FRAME_TYPE_TX_STATUS = 0x8B # FRAME_TYPE_TX_STATUS = 0x8B
FRAME_TYPE_RX_PACKET = 0x90 # FRAME_TYPE_RX_PACKET = 0x90
# ADDR16 選項 # # ADDR16 選項
DEST_ADDR16_UNICAST = b'\xFF\xFE' # DEST_ADDR16_UNICAST = b'\xFF\xFE'
DEST_ADDR16_BRAODCAST = b'\xFF\xFF' # 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): def __init__(self, at_handler: "ATCommandHandler" = None):
super().__init__() super().__init__()
self.at_handler = at_handler 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: def process_incoming(self, data: bytes) -> bytes:
@ -172,7 +173,11 @@ class XBeeFrameProcessor_Base(FrameProcessor):
def process_outgoing(self, data: bytes) -> bytes: def process_outgoing(self, data: bytes) -> bytes:
"""將數據封裝為 XBee API 傳輸幀""" """將數據封裝為 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: def _try_extract_frame(self) -> bytes:
@ -209,15 +214,15 @@ class XBeeFrameProcessor_Base(FrameProcessor):
""" """
frame_type = frame[3] 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] 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: if self.at_handler is not None:
self.at_handler.handle_frame(frame) self.at_handler.handle_frame(frame)
return None 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 # length = (frame[1] << 8) | frame[2] # for debug
# logger.info( # logger.info(
# f"TX Status raw={frame.hex()}, api_len={length}, " # f"TX Status raw={frame.hex()}, api_len={length}, "
@ -233,18 +238,18 @@ class XBeeFrameProcessor_Base(FrameProcessor):
@staticmethod @staticmethod
def _encapsulate( def _encapsulate(
data: bytes, data: bytes,
dest_addr64: bytes = DEST_ADDR64_BRAODCAST, dest_addr64: bytes,
dest_addr16: bytes = DEST_ADDR16_BRAODCAST, dest_addr16: bytes,
frame_id: int = 0x01, frame_id: int = 0x00,
options: int = 0x00,
) -> bytes: ) -> bytes:
""" """
payload 包成 XBee TX Request (0x10) payload 包成 XBee TX Request (0x10)
- 使用廣播地址 - 使用廣播地址
- 添加適當的頭部和校驗和 - 添加適當的頭部和校驗和
""" """
frame_type = XBeeFrameProcessor_Base.FRAME_TYPE_TX_REQUEST frame_type = ATCommandHandler.FRAME_TYPE_TX_REQUEST
broadcast_radius = 0x00 broadcast_radius = 0x00
options = 0x00
frame = struct.pack(">B", frame_type) + struct.pack(">B", frame_id) frame = struct.pack(">B", frame_type) + struct.pack(">B", frame_id)
frame += dest_addr64 + dest_addr16 frame += dest_addr64 + dest_addr16
@ -279,6 +284,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
這邊認定 esp_sysid 會等於 mavlink sysid 這邊認定 esp_sysid 會等於 mavlink sysid
poll time = 0.2305 * grant_bytes + 58.094 (R平方為0.979 高線性相關)
硬體產品系列 (Product Family) : XBP9B-DM 硬體產品系列 (Product Family) : XBP9B-DM
晶片世代: XBee PRO 900HP 200K 晶片世代: XBee PRO 900HP 200K
運作頻段: 900 MHz RF 運作頻段: 900 MHz RF
@ -304,15 +311,17 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
self.last_hello_time = last_hello_time self.last_hello_time = last_hello_time
self.remain_bytes = 0 # 剩餘 buffer 量 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.last_done_time = 0.0 # 最後送出Done的時間
self.received_len = 0 # 收到封包累計 self.received_len = 0 # 收到封包累計
self.poll_miss_count = 0 # poll 封包連續丟失的累計數
def __init__(self, at_handler: "ATCommandHandler" = None): def __init__(self, at_handler: "ATCommandHandler" = None):
super().__init__(at_handler) super().__init__(at_handler)
self.max_discovery_window_ms = 220 self.max_discovery_window_ms = 0.22
self.is_discovery_phase = False self.is_discovery_phase = False # 只是 show 狀態 用不太到
self.esp32_address_mapping = {} # bytes[Addr64] : Esp32DeviceInfo self.esp32_address_mapping = {} # bytes[Addr64] : Esp32DeviceInfo
self.operator_busy = False self.operator_busy = False
@ -326,21 +335,21 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
self.event_loop: Optional[asyncio.AbstractEventLoop] = None self.event_loop: Optional[asyncio.AbstractEventLoop] = None
self.serial_baudrate = 115200 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.poll_scheduler_state = pollStrategy.PollSchedulerState()
self.command_pending_event: Optional[asyncio.Event] = None self.command_pending_event: Optional[asyncio.Event] = None
self.poll_done_event: Optional[asyncio.Event] = None self.poll_done_event: Optional[asyncio.Event] = None
self.current_poll_address_64: Optional[bytes] = None self.current_poll_address_64: Optional[bytes] = None
self.current_poll_missing: bool = False
self.current_poll_time = 0 # 本次 poll 丟出的時間紀錄 用來運算傳輸效能 self.current_poll_time = 0 # 本次 poll 丟出的時間紀錄 用來運算傳輸效能
self.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間 self.last_discovery_time = 0.0 # 這個是最後做廣播 discovery 的時間
self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的 self.last_recieve_mavlink = 0.0 # 這個是最後收到 mavlink payload 時間 為了定義 poll-done 之間不要超時用的
self.mavPack_interval_timeout = 150 # mspoll 期間 MAVLink/DONE 最大閒置間隔 self.mavPack_interval_timeout = 150 # (msec) poll 期間 每筆 MAVLink 資料 /DONE 最大閒置間隔
self.discovery_interval_seconds = 200.0 # 每次做 discovery 程序的間隔時間 self.discovery_interval_seconds = 200.0 # (sec) 每次做 discovery 程序的間隔時間
self.device_offline_timeout = self.discovery_interval_seconds * 2 # 遠端沒有回應會被踢出 超時時限 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.guard_milliseconds = 50 # POLL DONE 的保底時間間隔
self.pending_manual_discovery = False 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.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: def set_writer(self, writer: Callable[[bytes], None]) -> None:
@ -361,6 +375,13 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
self.serial_baudrate = baudrate self.serial_baudrate = baudrate
def _ensure_async_primitives(self) -> None: def _ensure_async_primitives(self) -> None:
"""
因為 FrameProcessor 被歸屬在 SerialHandler 裡面
SerialHandler 物件在初始化時 是在 main thread
然而 以下兩個 Event 物件是在分支 thread 運作的
所以這裡不能放到 __init__
"""
if self.command_pending_event is None: if self.command_pending_event is None:
self.command_pending_event = asyncio.Event() self.command_pending_event = asyncio.Event()
if self.poll_done_event is None: if self.poll_done_event is None:
@ -371,12 +392,145 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
self._ensure_async_primitives() self._ensure_async_primitives()
self.command_pending_event.set() 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]: def _dispatch_frame(self, frame: bytes) -> Optional[bytes]:
frame_type = frame[3] 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) payload, sender_address_64 = self._decapsulate(frame)
if payload.startswith(self.HELLO_HEADER) and len(payload) == 5: 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() self.last_recieve_mavlink = time.time()
return payload 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: if self.at_handler is not None:
self.at_handler.handle_frame(frame) self.at_handler.handle_frame(frame)
return None return None
if frame_type == self.FRAME_TYPE_TX_STATUS: if frame_type == ATCommandHandler.FRAME_TYPE_TX_STATUS:
# length = (frame[1] << 8) | frame[2] self._handle_tx_status(frame)
# logger.debug( length = (frame[1] << 8) | frame[2] # for debug
# f"TX Status raw={frame.hex()}, api_len={length}, " logger.info(
# f"fid=0x{frame[4]:02X}, dest16=0x{(frame[5]<<8)|frame[6]:04X}, " f"TX Status raw={frame.hex()}, api_len={length}, "
# f"retry={frame[7]}, delivery={frame[8]}, discovery={frame[9]}" 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 return None
logger.warning(f"Unknown XBee frame type: 0x{frame_type:02X}") 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: 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) return self._encapsulate(
self.DISC_HEADER,
dest_addr16 = self.DEST_ADDR16_BRAODCAST,
dest_addr64 = self.DEST_ADDR64_BRAODCAST,
frame_id=0x05)
# 處理每個裝置回傳的 Hello 訊息 # 處理每個裝置回傳的 Hello 訊息
def handle_hello_report(self, payload: bytes, sender_address_64: bytes) -> None: def handle_hello_report(self, payload: bytes, sender_address_64: bytes) -> None:
@ -451,21 +610,18 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
# ---- POLL 封裝 / DONE 處理 ---- # ---- POLL 封裝 / DONE 處理 ----
def pack_poll(self, target_address_64: bytes, grant_bytes: int = 0) -> Optional[bytes]: def pack_poll(self, remote_device, grant_bytes: int = 0) -> Optional[bytes]:
# logger.debug(f"pack poll {target_address_64.hex()} / {grant_bytes}") # logger.debug(f"pack poll {remote_device.address_64.hex()} / {grant_bytes}")
remote_device = self.esp32_address_mapping.get(target_address_64)
if remote_device is None:
return None
poll_payload = self.POLL_HEADER + struct.pack( poll_payload = self.POLL_HEADER + struct.pack(
'>BH', remote_device.system_id, grant_bytes '>BH', remote_device.system_id, grant_bytes
) )
remote_device.received_len = 0
return self._encapsulate( return self._encapsulate(
poll_payload, poll_payload,
dest_addr64=remote_device.address_64, dest_addr64=remote_device.address_64,
dest_addr16=self.DEST_ADDR16_UNICAST, dest_addr16=self.DEST_ADDR16_UNICAST,
frame_id=0x02, frame_id=0x02,
options=0x01 # xbee auto retry
) )
def handle_done_report(self, payload: bytes, sender_address_64: bytes) -> None: 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 transmit_time = int((time.time() - self.current_poll_time)*1000) # ms
vehicle = vehicle_registry.get(system_id) 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: # 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 -= vehicle.rf_module.poll_data_sizes[0]
@ -500,77 +657,6 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
): ):
self.poll_done_event.set() 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 排程輔助 ----
# 計算下次要 poll 的對象跟大小 # 計算下次要 poll 的對象跟大小
@ -610,6 +696,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
def _should_run_discovery(self) -> bool: def _should_run_discovery(self) -> bool:
# 條件1. 目前沒有任何遠端ESP裝置被紀錄 或者 手動啟動 # 條件1. 目前沒有任何遠端ESP裝置被紀錄 或者 手動啟動
if (not self.esp32_address_mapping) or (self.pending_manual_discovery): 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) return (time.time() - self.last_discovery_time) > (1 + 5 * self.max_discovery_window_ms)
# 條件2. 每個固定週期 會做一次 # 條件2. 每個固定週期 會做一次
return (time.time() - self.last_discovery_time) > self.discovery_interval_seconds return (time.time() - self.last_discovery_time) > self.discovery_interval_seconds
@ -683,6 +770,8 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
}) })
return { return {
"operator_busy": self.operator_busy, "operator_busy": self.operator_busy,
"rf_busy": self.rf_busy,
"pending_tx_count": self._pending_tx_count,
"is_discovery_phase": self.is_discovery_phase, "is_discovery_phase": self.is_discovery_phase,
"gcs_queue_bytes": self.get_gcs_queue_byte_count(), "gcs_queue_bytes": self.get_gcs_queue_byte_count(),
"gcs_queue_packets": self.get_gcs_queue_packet_count(), "gcs_queue_packets": self.get_gcs_queue_packet_count(),
@ -746,16 +835,20 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
return return
# 4. (最優先) 處理從 UDP 過來的封包 並且透過 Serial 送出 # 4. (最優先) 處理從 UDP 過來的封包 並且透過 Serial 送出
if self.gcs_transmit_queue: if self.udp_transmit_queue:
await self.flush_gcs_transmit_queue() await self.flush_udp_transmit_queue()
return return
# 5. 檢測要不要跑 discovery 程序 # 5. RF 仍在傳輸,略過 discovery / poll
if self.rf_busy:
return
# 6. 檢測要不要跑 discovery 程序
if self._should_run_discovery(): if self._should_run_discovery():
await self._run_discovery() await self._run_discovery()
return return
# 6. POLL 程序 (若有手動先執行 若無自動則策略決策) # 7. POLL 程序 (若有手動先執行 若無自動則策略決策)
if self.pending_manual_poll is not None: if self.pending_manual_poll is not None:
target_system_id, grant_bytes = self.pending_manual_poll target_system_id, grant_bytes = self.pending_manual_poll
self.pending_manual_poll = None self.pending_manual_poll = None
@ -816,31 +909,40 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
# 這邊把要求封包最小量放在36 是讓載具端正好可以回傳一個有簽章的 HEARTBEAT 封包的大小 # 這邊把要求封包最小量放在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) remote_device = self.esp32_address_mapping.get(target_address_64)
if poll_frame is None: if remote_device is None:
return return None
# 打包 poll pack
poll_frame = self.pack_poll(remote_device, grant_bytes)
self._ensure_async_primitives() self._ensure_async_primitives()
self.operator_busy = True self.operator_busy = True
remote_device.received_len = 0
self.current_poll_address_64 = target_address_64 self.current_poll_address_64 = target_address_64
self.current_poll_time = time.time() self.current_poll_time = time.time()
self.current_poll_missing = False
self.poll_done_event.clear() self.poll_done_event.clear()
self.serial_writer(poll_frame) self.serial_writer(poll_frame)
self.last_recieve_mavlink = time.time() self.last_recieve_mavlink = time.time()
try: completed = await self._wait_poll_done_with_idle_timeout()
completed = await self._wait_poll_done_with_idle_timeout() if not completed:
if not completed: remote_device.remain_bytes -= remote_device.received_len
logger.warning( logger.warning(
f"POLL timeout address_64={target_address_64.hex()} " f"POLL timeout address_64={target_address_64.hex()} "
f"grant_bytes={grant_bytes}" f"grant_bytes={grant_bytes}"
) )
finally: # 判斷是否 poll pack 丟失 並紀錄連續丟失 清空 remain_bytes
self.operator_busy = False if self.current_poll_missing:
self.current_poll_address_64 = None remote_device.poll_miss_count += 1
await asyncio.sleep(self.guard_milliseconds / 1000.0) remote_device.remain_bytes = 0
self.operator_busy = False
self.current_poll_address_64 = None
# await asyncio.sleep(self.guard_milliseconds / 1000.0)
# poll 程序中 只要目標的載具還有在傳資訊 就不會終止資料的接收 # poll 程序中 只要目標的載具還有在傳資訊 就不會終止資料的接收
async def _wait_poll_done_with_idle_timeout(self) -> bool: 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(): while not self.poll_done_event.is_set():
if time.time() - self.last_recieve_mavlink >= idle_timeout_sec: if time.time() - self.last_recieve_mavlink >= idle_timeout_sec:
# logger.debug(f"mark A poll timeout")
return False return False
try: try:
await asyncio.wait_for( await asyncio.wait_for(
@ -861,7 +964,7 @@ class XBeeFrameProcessor_ESPv1(XBeeFrameProcessor_Base):
def stop_operator(self) -> None: def stop_operator(self) -> None:
self.operator_running = False self.operator_running = False
self.gcs_transmit_queue.clear() self.udp_transmit_queue.clear()
self.wake_operator() self.wake_operator()
@ -881,6 +984,12 @@ class ATCommandHandler:
""" """
FRAME_TYPE_AT_COMMAND = 0x08 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): def __init__(self, serial_port: str):
self.send_command_time = 0.0 # for dev self.send_command_time = 0.0 # for dev
@ -893,9 +1002,8 @@ class ATCommandHandler:
# 可擴展其他 AT 指令 # 可擴展其他 AT 指令
} }
# ---- 發送端 ---- # ---- 發送端 ---- # 這個不應該出現在這 越權了
def set_writer(self, writer): def set_writer(self, writer):
"""由 SerialHandler 注入實際寫入 serial 的 callable"""
self.writer = writer self.writer = writer
def send_command(self, request: ATRequest): def send_command(self, request: ATRequest):
@ -1083,15 +1191,9 @@ class UDPHandler(asyncio.DatagramProtocol):
logger.warning("Serial handler not set, dropping UDP packet") logger.warning("Serial handler not set, dropping UDP packet")
return 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) 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 ===================== # ====================== Serial Manager =====================
@ -1479,8 +1581,12 @@ if __name__ == '__main__':
processor = sm.get_espv1_processor(serial_id) processor = sm.get_espv1_processor(serial_id)
if processor is not None: if processor is not None:
processor.request_discovery() processor.request_discovery()
processor.request_poll(target_system_id=1) time.sleep(1)
processor.request_poll(target_system_id=1, grant_bytes=200) 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_status_snapshot())
print(processor.get_gcs_queue_byte_count()) print(processor.get_gcs_queue_byte_count())
@ -1494,7 +1600,7 @@ if __name__ == '__main__':
# sm.send_at_command(1, rssi_request) # sm.send_at_command(1, rssi_request)
# time.sleep(1) # time.sleep(1)
time.sleep(10) time.sleep(120)
sm.remove_serial_link(serial_id) sm.remove_serial_link(serial_id)
sm.shutdown() sm.shutdown()

Loading…
Cancel
Save