From 235ba61454e9b5fa1828fefdfe16472c222e25fe Mon Sep 17 00:00:00 2001 From: YikaiFu-cart Date: Thu, 13 Aug 2026 22:18:32 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=B7=BB=E5=8A=A0ACT=E5=8F=8C=E7=9B=B8?= =?UTF-8?q?=E6=9C=BA=E5=AE=9E=E6=97=B6=E9=A2=84=E8=A7=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../test/test_act_episode_recorder.py | 62 +++++- .../xr_rm_teleop/act_episode_recorder.py | 185 ++++++++++++++++++ 2 files changed, 245 insertions(+), 2 deletions(-) diff --git a/xr_rm_teleop/test/test_act_episode_recorder.py b/xr_rm_teleop/test/test_act_episode_recorder.py index bb864e4..18244f4 100644 --- a/xr_rm_teleop/test/test_act_episode_recorder.py +++ b/xr_rm_teleop/test/test_act_episode_recorder.py @@ -326,6 +326,22 @@ def test_sample_frame_metrics_separate_repeats_skips_and_regressions(): assert regression == 0 +def test_camera_buffer_reports_two_second_rolling_fps(): + buffer = CameraBuffer(maxlen=4) + start_ns = 10_000_000_000 + for index in range(61): + buffer.push( + CameraFrame( + _image(index), + index, + float(index), + start_ns + index * 33_333_333, + ) + ) + + assert buffer.stats().rolling_fps == pytest.approx(30.0, rel=0.02) + + requires_h5py = pytest.mark.skipif( h5py is None, reason="h5py is not installed", @@ -763,6 +779,48 @@ def _recorder_for_test(tmp_path): return recorder +def test_preview_failure_does_not_change_recording_state(): + recorder = object.__new__(ActEpisodeRecorder) + recorder._session = _recording_session() + recorder._preview_stop = threading.Event() + recorder._preview_thread = None + recorder._logger = _Logger() + recorder.get_logger = lambda: recorder._logger + + class FailingCv2: + WINDOW_NORMAL = 0 + + @staticmethod + def namedWindow(*_args): + raise RuntimeError("no display") + + recorder._preview_loop(FailingCv2()) + + assert recorder.state is RecordingState.RECORDING + assert recorder._logger.warnings == [ + "ACT双相机预览已停用:no display" + ] + + +def test_close_stops_preview_before_cameras_and_releases_lock(): + events = [] + recorder = object.__new__(ActEpisodeRecorder) + recorder._stop_preview = lambda: events.append("preview") + recorder._high_camera = SimpleNamespace( + stop=lambda: events.append("high") + ) + recorder._wrist_camera = SimpleNamespace( + stop=lambda: events.append("wrist") + ) + recorder._directory_lock = SimpleNamespace( + release=lambda: events.append("lock") + ) + + recorder.close() + + assert events == ["preview", "high", "wrist", "lock"] + + @requires_h5py def test_preflight_requires_open_gripper_fresh_inputs_and_disk_space(tmp_path): recorder = _recorder_for_test(tmp_path) @@ -791,8 +849,8 @@ def _push_recording_frames(recorder, control_ns, frame_number): recorder._wrist_camera.buffer.push( CameraFrame( _image(3), - frame_number, - float(frame_number), + frame_number + 1000, + float(frame_number + 1000), control_ns - 3_000_000, ) ) diff --git a/xr_rm_teleop/xr_rm_teleop/act_episode_recorder.py b/xr_rm_teleop/xr_rm_teleop/act_episode_recorder.py index 7ffd417..b36c7a9 100644 --- a/xr_rm_teleop/xr_rm_teleop/act_episode_recorder.py +++ b/xr_rm_teleop/xr_rm_teleop/act_episode_recorder.py @@ -638,6 +638,7 @@ class CameraStats: frame_number_regression_count: int first_host_monotonic_ns: int | None last_host_monotonic_ns: int | None + recent_host_monotonic_ns: tuple[int, ...] @property def drop_ratio(self) -> float: @@ -655,6 +656,20 @@ class CameraStats: return 0.0 return (self.frame_count - 1) * 1e9 / elapsed_ns + @property + def rolling_fps(self) -> float: + if len(self.recent_host_monotonic_ns) < 2: + return 0.0 + elapsed_ns = ( + self.recent_host_monotonic_ns[-1] + - self.recent_host_monotonic_ns[0] + ) + if elapsed_ns <= 0: + return 0.0 + return ( + (len(self.recent_host_monotonic_ns) - 1) * 1e9 / elapsed_ns + ) + class CameraBuffer: def __init__(self, *, maxlen: int = 4) -> None: @@ -668,6 +683,7 @@ class CameraBuffer: self._first_host_monotonic_ns: int | None = None self._last_host_monotonic_ns: int | None = None self._last_frame_number: int | None = None + self._recent_host_monotonic_ns: deque[int] = deque() def push(self, frame: CameraFrame) -> None: with self._lock: @@ -683,6 +699,13 @@ class CameraBuffer: if self._first_host_monotonic_ns is None: self._first_host_monotonic_ns = frame.host_monotonic_ns self._last_host_monotonic_ns = frame.host_monotonic_ns + self._recent_host_monotonic_ns.append(frame.host_monotonic_ns) + cutoff_ns = frame.host_monotonic_ns - 2_000_000_000 + while ( + self._recent_host_monotonic_ns + and self._recent_host_monotonic_ns[0] < cutoff_ns + ): + self._recent_host_monotonic_ns.popleft() self._frames.append(frame) def snapshot(self) -> tuple[CameraFrame, ...]: @@ -699,6 +722,9 @@ class CameraBuffer: ), first_host_monotonic_ns=self._first_host_monotonic_ns, last_host_monotonic_ns=self._last_host_monotonic_ns, + recent_host_monotonic_ns=tuple( + self._recent_host_monotonic_ns + ), ) @@ -1049,6 +1075,8 @@ def _twist_values(twist: Any) -> np.ndarray: class ActEpisodeRecorder(Node): + PREVIEW_WINDOW = "ACT - D455 Global / D405 Right Wrist" + def __init__(self) -> None: super().__init__("act_episode_recorder") defaults = { @@ -1147,6 +1175,8 @@ class ActEpisodeRecorder(Node): self._episode_index: int | None = None self._camera_baselines: tuple[CameraStats, CameraStats] | None = None self._saving_deadline_ns: int | None = None + self._preview_stop = threading.Event() + self._preview_thread: threading.Thread | None = None self._high_camera = RealSenseCamera( str(parameters["cam_high_serial"]), @@ -1163,6 +1193,8 @@ class ActEpisodeRecorder(Node): except Exception as exc: self._camera_start_error = str(exc) self.get_logger().error(f"ACT相机启动失败:{exc}") + if self._camera_start_error is None: + self._start_preview() self._status_pub = self.create_publisher( String, @@ -1226,6 +1258,158 @@ class ActEpisodeRecorder(Node): def state(self) -> RecordingState: return self._session.state + @staticmethod + def _preview_tile( + role: str, + frame: CameraFrame | None, + stats: CameraStats, + cv2: Any, + ) -> np.ndarray: + if frame is None: + tile = np.zeros((480, 640, 3), dtype=np.uint8) + frame_number = "-" + else: + tile = cv2.cvtColor(frame.image, cv2.COLOR_RGB2BGR) + frame_number = str(frame.frame_number) + tile = tile.copy() + cv2.rectangle(tile, (0, 0), (640, 74), (0, 0, 0), -1) + lines = ( + f"{role} FPS {stats.rolling_fps:.1f} Frame {frame_number}", + f"Received {stats.frame_count} Dropped " + f"{stats.dropped_frames} ({stats.drop_ratio:.2%})", + ) + for index, line in enumerate(lines): + cv2.putText( + tile, + line, + (10, 28 + index * 30), + cv2.FONT_HERSHEY_SIMPLEX, + 0.65, + (255, 255, 255), + 1, + cv2.LINE_AA, + ) + return tile + + def _compose_preview(self, cv2: Any) -> np.ndarray: + high_frames = self._high_camera.buffer.snapshot() + wrist_frames = self._wrist_camera.buffer.snapshot() + high = high_frames[-1] if high_frames else None + wrist = wrist_frames[-1] if wrist_frames else None + image = np.hstack( + ( + self._preview_tile( + "GLOBAL D455", + high, + self._high_camera.buffer.stats(), + cv2, + ), + self._preview_tile( + "RIGHT WRIST D405", + wrist, + self._wrist_camera.buffer.stats(), + cv2, + ), + ) + ) + + now_ns = self._now_ns() + high_age = ( + f"{(now_ns - high.host_monotonic_ns) * 1e-6:.1f} ms" + if high is not None + else "-" + ) + wrist_age = ( + f"{(now_ns - wrist.host_monotonic_ns) * 1e-6:.1f} ms" + if wrist is not None + else "-" + ) + skew = ( + f"{abs(high.host_monotonic_ns - wrist.host_monotonic_ns) * 1e-6:.1f} ms" + if high is not None and wrist is not None + else "-" + ) + episode_index = self._episode_index + episode = ( + f"episode_{episode_index}" + if episode_index is not None + else "-" + ) + store = self._store + samples = store.count if store is not None else 0 + footer = np.zeros((80, image.shape[1], 3), dtype=np.uint8) + lines = ( + f"Age high={high_age} wrist={wrist_age} Camera skew={skew}", + f"ACT {self.state.value} {episode} Samples {samples}", + ) + for index, line in enumerate(lines): + cv2.putText( + footer, + line, + (10, 30 + index * 32), + cv2.FONT_HERSHEY_SIMPLEX, + 0.7, + (255, 255, 255), + 1, + cv2.LINE_AA, + ) + return np.vstack((image, footer)) + + def _start_preview(self) -> None: + if not ( + os.environ.get("DISPLAY") or os.environ.get("WAYLAND_DISPLAY") + ): + self.get_logger().warn( + "未检测到桌面显示环境,ACT双相机预览已停用。" + ) + return + try: + import cv2 + except ImportError as exc: + self.get_logger().warn( + f"OpenCV不可用,ACT双相机预览已停用:{exc}" + ) + return + self._preview_stop.clear() + self._preview_thread = threading.Thread( + target=self._preview_loop, + args=(cv2,), + name="act_camera_preview", + daemon=True, + ) + self._preview_thread.start() + + def _preview_loop(self, cv2: Any) -> None: + try: + cv2.namedWindow(self.PREVIEW_WINDOW, cv2.WINDOW_NORMAL) + while not self._preview_stop.is_set(): + cv2.imshow( + self.PREVIEW_WINDOW, + self._compose_preview(cv2), + ) + key = cv2.waitKey(1) & 0xFF + if key in (ord("q"), ord("Q"), 27): + break + if cv2.getWindowProperty( + self.PREVIEW_WINDOW, + cv2.WND_PROP_VISIBLE, + ) < 1: + break + self._preview_stop.wait(0.1) + except Exception as exc: + self.get_logger().warn(f"ACT双相机预览已停用:{exc}") + finally: + try: + cv2.destroyWindow(self.PREVIEW_WINDOW) + except Exception: + pass + + def _stop_preview(self) -> None: + self._preview_stop.set() + if self._preview_thread is not None: + self._preview_thread.join(timeout=2.0) + self._preview_thread = None + def _publish_state( self, state: RecordingState, @@ -1729,6 +1913,7 @@ class ActEpisodeRecorder(Node): self._reject_current(reason, interrupted=True) def close(self) -> None: + self._stop_preview() self._high_camera.stop() self._wrist_camera.stop() self._directory_lock.release()