tmp commit

chiyu
Chiyu Chen 4 weeks ago
parent 82d5e90460
commit 2542663e9e

@ -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) # 其他自定義狀態
@ -276,6 +277,11 @@ class RFModule:
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"""
self.status.rssi = rssi

@ -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 之間不要超時用的
@ -347,6 +346,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)
# ---- 注入與設定 ----
def set_writer(self, writer: Callable[[bytes], None]) -> None:
@ -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()
if vehicle:=vehicle_registry.get(system_id):
if not vehicle.rf_module:
vehicle.rf_module = RFModule(RFModuleType.XBEE_ESP)
# 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()}"
# )
# ---- 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)
# 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,
)
# 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()

Loading…
Cancel
Save