diff --git a/README.md b/README.md index 1325297..b4e176d 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,8 @@ TOML 解析检查。若还要核对外置 OmniSocketGo 版本,可设置任务 ## 当前按键 - 左 Z + 右 C 连续 3 秒:开始遥操;再次连续 3 秒:结束并限速回 Home。 + 若 Hub 尚未运行,本次长按只检测一次并立即结束,不发送控制数据;服务器启动后须先 + 稳定松开,再重新长按 3 秒。pending 和活动会话都不会自动重试或自动恢复。 - 右 C + 左摇杆上下:HBWALK 前进/后退。 - 左 Z + 右摇杆左右:HBWALK 原地转向。 - 右 B 连续 3 秒:右手进入厂商“单食指”姿态;松开至少 0.5 秒后再次连续 diff --git a/docs/天工3.0本地同构臂遥操迁移部署指南.md b/docs/天工3.0本地同构臂遥操迁移部署指南.md index cd53ce1..e90a7d9 100644 --- a/docs/天工3.0本地同构臂遥操迁移部署指南.md +++ b/docs/天工3.0本地同构臂遥操迁移部署指南.md @@ -427,6 +427,15 @@ omnisocket_server = "203.0.113.10:15000" 不需要。`127.0.0.1` 永远表示 EAI 本机,和网卡地址无关。 +### Hub 没启动时已经长按了 Z+C,服务器恢复后怎么办? + +本次长按只做一次后台连接检测。Hub 未启动、机器人 Peer 不存在或首帧发送失败时, +`teleop_start_pending` 会回到 `false`;Hub 离线时没有业务帧发出,Peer/发送错误时机器人 +不会收到控制数据。旧请求不会在服务器稍后恢复 +时自动重试。稳定松开 Z+C 至少 0.5 秒,待 Hub 与机器人接收端就绪后重新长按 3 秒即可。 +后台检测期间 sender 仍持续读取 5003,因此松键不会因连接等待而漏掉。活动会话断网同样 +按故障关闭处理,必须稳定松开后重新长按。 + ### 只改机器人 Peer,不改 EAI 可以吗? 不可以。EAI 的 `--target-peer` 必须等于机器人的 `omnisocket_peer_id`,机器人 diff --git a/tg3_omnisocket_transport/README.md b/tg3_omnisocket_transport/README.md index 0bf8095..a63ef8c 100644 --- a/tg3_omnisocket_transport/README.md +++ b/tg3_omnisocket_transport/README.md @@ -36,10 +36,19 @@ xTELE 已处理的标量或六维 BrainCoRevo2 目标。5002 的界面事件和 xTELE 帧;连续长按左 Z + 右 C 3 秒后才生成新的 128-bit `session_id` 并开始发送。 再次长按 3 秒时发送最后一帧 `stop`,等待有界 KCP 刷新后关闭 OmniSocket Session; 之后 `frames_sent/bytes_sent` 不再增长,待机既没有 xTELE 业务帧,也没有该 sender 的 -底层注册/心跳。即使 Hub 离线,sender 仍持续读取本机 xTELE 并保持会话关闭;启动连接 -失败会使本次会话失效,必须松开组合键后重新长按,不会在 Hub 恢复时续发旧动作。 +底层注册/心跳。即使 Hub 离线,sender 仍持续读取本机 xTELE 并保持会话关闭。Z+C 满 +3 秒后先进入 `START_PENDING`,后台连接不会阻塞本机按键采样;本次长按只尝试连接一次。 +若 Hub 未启动,本次连接失败且不会发出业务帧;若机器人 Peer 不存在或首帧发送失败, +机器人不会收到控制数据且 sender 在收到错误后立即关闭。这次请求都必须稳定松开后再做 +一次新的 3 秒长按;Hub 稍后恢复不会让旧请求自行启动。 +任一键松开、按键数据畸形或 5003 失联也会立即取消 pending。连接成功后还必须看到两帧 +时间戳最终递增且仍按住 Z+C 的新鲜 5003 数据,才发送第一帧 START,连接前排队的旧 +动作不会被续发。 活动期间若 KCP 反馈超过 `500 ms` 未更新,或 `snd_queue+snd_buffer` 超过 100 帧,sender 立即作废会话并关闭 Session,从源头丢弃待发队列,防止网络恢复后回放旧动作。 +5003 的源时间戳在活动期间也必须持续推进;允许相邻毫秒值短暂重复,但超过 `250 ms` +不再推进、时间戳倒退、持续畸形 JSON 或超长帧都会立即作废会话。操作员 STOP 无论最终 +业务帧能否编码或发送,EAI 都会在本地关闭 Session,机器人端由输入超时保护停止并回 Home。 每个业务 JSON 的 `tg3_transport` 由 sender 强制覆盖,不能由 5003 输入伪造: @@ -55,7 +64,8 @@ xTELE 帧;连续长按左 Z + 右 C 3 秒后才生成新的 128-bit `session_i 当前用户服务将 `start` 连续发送 500 个源帧,避免接收端只取最新帧时错过唯一启动事件;机器人对同一 `session_id` 只做一次启动安全检查。`stop` 是该会话最后一个业务帧。服务重启、源数据 -失联或网络错误都会使会话失效,恢复后必须先松开组合键,再重新长按 3 秒。 +失联或网络错误都会使会话失效,恢复后必须先松开组合键,再重新长按 3 秒;pending 和 +活动会话都不会自动重试或自动恢复。 `tg3_local_teleop` 内部直接拒绝非预期发送端、 乱序、格式错误以及相对本次连接最低包龄额外排队超过 300 ms 的数据;相对包龄会消除 两端系统时钟的固定偏差。遥操桥仍有 250 ms 输入失联保护。 diff --git a/tg3_omnisocket_transport/omnisocket_xtele_sender.py b/tg3_omnisocket_transport/omnisocket_xtele_sender.py index b81ab53..586abfb 100755 --- a/tg3_omnisocket_transport/omnisocket_xtele_sender.py +++ b/tg3_omnisocket_transport/omnisocket_xtele_sender.py @@ -20,6 +20,7 @@ import os from pathlib import Path import signal import struct +import threading import time from typing import Any import uuid @@ -56,6 +57,8 @@ class XteleSender: self.last_status_at = 0.0 self.last_frame_at = 0.0 self.last_source_frame_at = 0.0 + self.last_source_timestamp: float | None = None + self.last_source_timestamp_advanced_at = 0.0 self.last_command_at = 0.0 self.last_error = "" self.teleop_active = False @@ -67,6 +70,21 @@ class XteleSender: self.combo_release_started_at: float | None = None # A service restart must never turn an already-held combo into a start. self.require_combo_release = True + # START is a two-phase transaction. A physical 3-second hold creates + # a pending request; only a successful Hub connection followed by a + # newer raw frame that still shows Z+C pressed may send START. + self.start_pending = False + self.start_wait_fresh_after_connect = False + self.start_fresh_barrier_timestamp: float | None = None + self.latest_combo_state: bool | None = None + self.last_start_failure = "" + # Session.connect() may wait for the Hub registration timeout. Run it + # outside the input loop so a release or stale xTELE source is still + # observed immediately while an attempt is in flight. + self.connect_attempt_generation = 0 + self.connect_thread: threading.Thread | None = None + self.connect_result_lock = threading.Lock() + self.connect_result: tuple[int, Session | None, str] | None = None # Keep xTELE's processed right-hand stream isolated for the complete # right-B press and stable-release transaction. The robot runs the # authoritative three-second gesture toggle from the raw B state. @@ -92,7 +110,15 @@ class XteleSender: "frames_suppressed_inactive": 0, "teleop_starts": 0, "teleop_stops": 0, + "teleop_stop_send_failures": 0, "teleop_aborts": 0, + "teleop_start_requests": 0, + "start_connect_attempts": 0, + "start_connect_failures": 0, + "start_send_failures": 0, + "start_pending_cancels": 0, + "source_timestamp_errors": 0, + "source_timestamp_resets": 0, } signal.signal(signal.SIGINT, self._request_stop) signal.signal(signal.SIGTERM, self._request_stop) @@ -100,28 +126,6 @@ class XteleSender: def _request_stop(self, _signum: int, _frame: Any) -> None: self.stop = True - def connect(self) -> bool: - self.close_session() - session = Session() - try: - session.connect( - server_addr=self.args.server, - peer_id=self.args.peer_id, - **CONTROL_DEFAULTS, - ) - except OSError as exc: - self.last_error = f"OmniSocket connect failed: {exc}" - session.close() - self.write_status(force=True) - return False - self.session = session - self.session_connected_at = time.monotonic() - self.counters["connected"] = 1 - self.counters["reconnects"] += 1 - self.last_error = "" - self.write_status(force=True) - return True - def close_session(self) -> None: if self.session is not None: try: @@ -161,6 +165,12 @@ class XteleSender: "last_source_frame_age_s": None if self.last_source_frame_at == 0.0 else round(time.monotonic() - self.last_source_frame_at, 4), + "last_source_timestamp": self.last_source_timestamp, + "source_timestamp_advance_age_s": None + if self.last_source_timestamp_advanced_at == 0.0 + else round( + time.monotonic() - self.last_source_timestamp_advanced_at, 4 + ), "teleop_active": self.teleop_active, "teleop_session_id": self.teleop_session_id, "teleop_session_seq": self.teleop_session_seq, @@ -171,6 +181,13 @@ class XteleSender: "teleop_combo_release_hold_s": 0.0 if self.combo_release_started_at is None else round(time.monotonic() - self.combo_release_started_at, 2), + "teleop_start_pending": self.start_pending, + "teleop_start_waiting_fresh_frame": ( + self.start_wait_fresh_after_connect + ), + "teleop_start_connect_inflight": self._connect_inflight(), + "teleop_latest_combo_state": self.latest_combo_state, + "last_start_failure": self.last_start_failure, "right_b_processed_merge_suppressed": self.right_b_merge_suppressed, "right_b_merge_release_hold_s": 0.0 if self.right_b_release_started_at is None @@ -197,6 +214,106 @@ class XteleSender: pass self.last_status_at = now + def _connect_inflight(self) -> bool: + thread = self.connect_thread + return thread is not None and thread.is_alive() + + def _connect_worker(self, generation: int) -> None: + session: Session | None = None + failure = "" + try: + session = Session() + session.connect( + server_addr=self.args.server, + peer_id=self.args.peer_id, + **CONTROL_DEFAULTS, + ) + except OSError as exc: + failure = f"OmniSocket connect failed: {exc}" + if session is not None: + try: + session.close() + except OSError: + pass + session = None + except Exception as exc: + failure = f"unexpected OmniSocket connect failure: {exc}" + if session is not None: + try: + session.close() + except OSError: + pass + session = None + with self.connect_result_lock: + self.connect_result = (generation, session, failure) + + def _start_connect_attempt(self) -> bool: + if self._connect_inflight(): + return False + with self.connect_result_lock: + # A worker can publish its result just after the main loop polls + # it. Let the next poll adopt/discard that result instead of + # overwriting (and leaking) a connected Session. + if self.connect_result is not None: + return False + self.connect_attempt_generation += 1 + generation = self.connect_attempt_generation + thread = threading.Thread( + target=self._connect_worker, + args=(generation,), + name="tg3-omnisocket-start-connect", + daemon=True, + ) + self.connect_thread = thread + thread.start() + return True + + def _invalidate_connect_attempt(self) -> None: + # The SDK connect call is not cancellable. Changing the generation + # makes its eventual result unusable; the main loop will close it. + self.connect_attempt_generation += 1 + + def _poll_connect_result(self, now: float) -> str | None: + with self.connect_result_lock: + result = self.connect_result + self.connect_result = None + if result is None: + return None + generation, session, failure = result + thread = self.connect_thread + if thread is not None: + thread.join(timeout=0.0) + self.connect_thread = None + if generation != self.connect_attempt_generation or not self.start_pending: + if session is not None: + try: + session.close() + except OSError: + pass + return "discarded" + if session is None: + self.counters["start_connect_failures"] += 1 + self.last_error = failure or "OmniSocket Hub connection failed" + # One physical hold makes exactly one connection attempt. A Hub + # failure consumes no robot command, but it does require a stable + # release before the operator may make a new three-second hold. + self._cancel_pending_start(self.last_error) + self.write_status(force=True) + return "failed" + + self.session = session + self.session_connected_at = now + self.counters["connected"] = 1 + self.counters["reconnects"] += 1 + self.last_error = "" + # Do not trust anything queued while registration was in progress. + # Two advancing raw timestamps after adoption prove the source is + # still live before START can leave EAI. + self.start_wait_fresh_after_connect = True + self.start_fresh_barrier_timestamp = None + self.write_status(force=True) + return "connected" + def _drain_responses(self) -> None: if self.session is None: return @@ -478,24 +595,38 @@ class XteleSender: return encoded, bool(merged_sides) @staticmethod - def _start_stop_pressed(data: dict[str, object]) -> bool: + def _start_stop_pressed(data: dict[str, object]) -> bool | None: + """Return explicit Z+C state, or ``None`` for malformed buttons.""" + try: buttons = data["button"] left = buttons["left"] # type: ignore[index] right = buttons["right"] # type: ignore[index] - return ( - len(left) >= 3 # type: ignore[arg-type] - and len(right) >= 3 # type: ignore[arg-type] - and bool(left[2]) # type: ignore[index] - and bool(right[2]) # type: ignore[index] - ) - except (KeyError, TypeError): - return False + if len(left) < 3 or len(right) < 3: # type: ignore[arg-type] + return None + left_z = left[2] # type: ignore[index] + right_c = right[2] # type: ignore[index] + except (IndexError, KeyError, TypeError): + return None + values: list[bool] = [] + for value in (left_z, right_c): + if isinstance(value, bool): + values.append(value) + elif isinstance(value, int) and value in (0, 1): + values.append(bool(value)) + else: + return None + return all(values) def _update_teleop_gate( self, now: float, data: dict[str, object] ) -> str | None: pressed = self._start_stop_pressed(data) + self.latest_combo_state = pressed + if pressed is None: + self.combo_started_at = None + self.combo_release_started_at = None + return None if self.require_combo_release: self.combo_started_at = None if pressed: @@ -528,33 +659,149 @@ class XteleSender: self.counters["teleop_stops"] += 1 return "stop" - self.teleop_active = True + self.start_pending = True self.teleop_session_id = uuid.uuid4().hex self.teleop_session_seq = 0 self.start_markers_remaining = self.args.start_marker_frames - self.counters["teleop_starts"] += 1 - return "start" + self.start_wait_fresh_after_connect = False + self.last_start_failure = "" + self.last_error = "" + self.counters["teleop_start_requests"] += 1 + return "start_pending" - def _abort_teleop(self, reason: str) -> None: - if self.teleop_active: - self.counters["teleop_aborts"] += 1 + def _clear_pending_start_state(self) -> None: + self.start_pending = False + self.start_wait_fresh_after_connect = False + self.start_fresh_barrier_timestamp = None + + def _cancel_pending_start(self, reason: str) -> None: + if not self.start_pending: + return + self._invalidate_connect_attempt() + self.close_session() + self._clear_pending_start_state() self.teleop_active = False self.teleop_session_id = None self.teleop_session_seq = 0 self.start_markers_remaining = 0 self.combo_started_at = None self.require_combo_release = True - self.last_error = reason + self.last_start_failure = reason + self.counters["start_pending_cancels"] += 1 - def _send_payload(self, payload: bytes) -> bool: - if self.session is None: - self._abort_teleop("OmniSocket session is unavailable") + def _pending_connect_ready(self, _now: float) -> bool: + if not self.start_pending or self.latest_combo_state is not True: return False + if self.session is not None: + return False + if self._connect_inflight(): + return False + return True + + def _commit_pending_start_sent(self, _now: float) -> None: + self.counters["teleop_starts"] += 1 + self.start_pending = False + self.start_wait_fresh_after_connect = False + self.start_fresh_barrier_timestamp = None + self.teleop_active = True + self.last_start_failure = "" + + def _abort_teleop(self, reason: str) -> None: + if self.teleop_active or self.start_pending: + self.counters["teleop_aborts"] += 1 + self._invalidate_connect_attempt() + self.teleop_active = False + self._clear_pending_start_state() + self.teleop_session_id = None + self.teleop_session_seq = 0 + self.start_markers_remaining = 0 + self.combo_started_at = None + self.require_combo_release = True + self.last_error = reason + self.last_start_failure = reason + + @staticmethod + def _raw_source_timestamp(data: dict[str, object]) -> float | None: + value = data.get("timestamp") + if isinstance(value, bool) or not isinstance(value, (int, float)): + return None + timestamp = float(value) + return timestamp if math.isfinite(timestamp) else None + + def _check_source_clock( + self, now: float, data: dict[str, object] + ) -> tuple[bool, str]: + """Reject a missing, regressing or persistently frozen xTELE clock.""" + + timestamp = self._raw_source_timestamp(data) + if timestamp is None: + return False, "local xTELE timestamp is missing or malformed" + previous = self.last_source_timestamp + if previous is None or timestamp > previous: + self.last_source_timestamp = timestamp + self.last_source_timestamp_advanced_at = now + return True, "" + if timestamp == previous: + if self.last_source_timestamp_advanced_at == 0.0: + self.last_source_timestamp_advanced_at = now + return True, "" + if ( + now - self.last_source_timestamp_advanced_at + <= self.args.source_timeout_s + ): + return True, "" + return False, "local xTELE timestamp stopped advancing" + + if not self.teleop_active and not self.start_pending: + # A producer restart while idle is harmless only after a fresh + # release. Reset the clock baseline but never inherit a held + # Z+C combination across that restart. + self.last_source_timestamp = timestamp + self.last_source_timestamp_advanced_at = now + self.combo_started_at = None + self.combo_release_started_at = None + self.require_combo_release = True + self.latest_combo_state = None + self.counters["source_timestamp_resets"] += 1 + return True, "" + return False, "local xTELE timestamp moved backwards" + + def _check_pending_start_freshness( + self, data: dict[str, object] + ) -> tuple[str, str]: + """Gate START on two advancing frames received after connect adoption.""" + + if not self.start_wait_fresh_after_connect: + return "ready", "" + source_timestamp = self._raw_source_timestamp(data) + if source_timestamp is None: + return ( + "cancel", + "pending START cannot prove a fresh xTELE timestamp", + ) + if self.start_fresh_barrier_timestamp is None: + self.start_fresh_barrier_timestamp = source_timestamp + return "wait", "" + if source_timestamp < self.start_fresh_barrier_timestamp: + return ( + "cancel", + "pending START xTELE timestamp moved backwards", + ) + if source_timestamp == self.start_fresh_barrier_timestamp: + # xTELE runs near 100 Hz while its millisecond timestamp can repeat + # for adjacent frames. Keep waiting; never treat a duplicate as + # proof of freshness and never force a new physical button cycle. + return "wait", "" + self.start_wait_fresh_after_connect = False + self.start_fresh_barrier_timestamp = None + return "ready", "" + + def _try_send_payload(self, payload: bytes) -> tuple[bool, str]: + if self.session is None: + return False, "OmniSocket session is unavailable" unhealthy = self._session_unhealthy_reason() if unhealthy: - self._abort_teleop(unhealthy) - self.close_session() - return False + return False, unhealthy self.packet_sequence += 1 packet = HEADER.pack( MAGIC, self.packet_sequence, time.time_ns(), len(payload) @@ -562,14 +809,20 @@ class XteleSender: try: self.session.send(to=self.args.target_peer, data=packet) except OSError as exc: - self._abort_teleop(f"OmniSocket send failed: {exc}") - self.close_session() - return False + return False, f"OmniSocket send failed: {exc}" self.counters["frames_sent"] += 1 self.counters["bytes_sent"] += len(payload) self.last_frame_at = time.monotonic() self.last_error = "" - return True + return True, "" + + def _send_payload(self, payload: bytes) -> bool: + sent, reason = self._try_send_payload(payload) + if sent: + return True + self._abort_teleop(reason) + self.close_session() + return False def _session_unhealthy_reason(self) -> str | None: if self.session is None: @@ -620,6 +873,16 @@ class XteleSender: self._drain_responses() time.sleep(0.01) + def _finish_stop_session(self) -> None: + """Close an operator STOP locally even if its final frame failed.""" + + self.teleop_active = False + self.teleop_session_id = None + self.teleop_session_seq = 0 + self.start_markers_remaining = 0 + self._flush_session() + self.close_session() + def run(self) -> int: context = zmq.Context() source = context.socket(zmq.SUB) @@ -643,6 +906,10 @@ class XteleSender: try: while not self.stop: events = dict(poller.poll(100)) + # Connection attempts run in a worker so this loop can keep + # observing physical releases and source freshness. Only the + # main thread ever adopts a completed Session. + self._poll_connect_result(time.monotonic()) if command_source is not None and command_source in events: command_raw = command_source.recv() self.counters["command_frames_received"] += 1 @@ -669,6 +936,10 @@ class XteleSender: ): self.combo_started_at = None self.combo_release_started_at = None + if self.start_pending: + self._cancel_pending_start( + "local xTELE source became stale during START" + ) if self.teleop_active: self._abort_teleop( "local xTELE source became stale; a new Z+C hold " @@ -683,6 +954,17 @@ class XteleSender: self.counters["frames_received"] += 1 self.counters["bytes_received"] += len(raw) if not raw or len(raw) > MAX_PAYLOAD_BYTES: + self.combo_started_at = None + self.combo_release_started_at = None + if self.start_pending: + self._cancel_pending_start( + "malformed local xTELE frame during START" + ) + if self.teleop_active: + self._abort_teleop( + "malformed local xTELE frame during active session" + ) + self.close_session() self.counters["dropped_malformed"] += 1 continue @@ -695,11 +977,33 @@ class XteleSender: # physical three-second start/stop hold. self.combo_started_at = None self.combo_release_started_at = None + if self.start_pending: + self._cancel_pending_start( + "malformed local xTELE JSON during START" + ) + if self.teleop_active: + self._abort_teleop( + "malformed local xTELE JSON during active session" + ) + self.close_session() self.counters["dropped_malformed"] += 1 continue now = time.monotonic() self.last_source_frame_at = now + clock_ok, clock_reason = self._check_source_clock(now, data) + if not clock_ok: + self.combo_started_at = None + self.combo_release_started_at = None + self.latest_combo_state = None + self.counters["source_timestamp_errors"] += 1 + self.counters["dropped_malformed"] += 1 + if self.start_pending: + self._cancel_pending_start(clock_reason) + if self.teleop_active: + self._abort_teleop(clock_reason) + self.close_session() + continue transition = self._update_teleop_gate(now, data) command = None @@ -723,6 +1027,86 @@ class XteleSender: now, data, processed_right ) + if self.start_pending and self.latest_combo_state is not True: + if self.latest_combo_state is False: + pending_cancel_reason = ( + "Z+C released before START reached the robot" + ) + else: + pending_cancel_reason = ( + "malformed Z+C state during pending START" + ) + self._cancel_pending_start(pending_cancel_reason) + + if self.start_pending: + if self.session is None: + if not self._pending_connect_ready(now): + self.counters["frames_suppressed_inactive"] += 1 + self._drain_responses() + self.write_status() + continue + if self._start_connect_attempt(): + self.counters["start_connect_attempts"] += 1 + self.counters["frames_suppressed_inactive"] += 1 + self.write_status(force=True) + continue + + freshness, freshness_reason = ( + self._check_pending_start_freshness(data) + ) + if freshness == "wait": + self.counters["frames_suppressed_inactive"] += 1 + self.write_status() + continue + if freshness == "cancel": + self._cancel_pending_start(freshness_reason) + self.counters["frames_suppressed_inactive"] += 1 + continue + + session_id = self.teleop_session_id + if not session_id: + self._cancel_pending_start( + "pending START session ID is unavailable" + ) + self.counters["frames_suppressed_inactive"] += 1 + continue + self.teleop_session_seq += 1 + try: + payload, merged = self._build_payload( + data, + command, + session_id, + self.teleop_session_seq, + "start", + "", + suppress_processed_right, + ) + except ValueError: + self._cancel_pending_start( + "cannot encode a fresh pending START frame" + ) + self.counters["dropped_malformed"] += 1 + continue + if len(payload) > MAX_PAYLOAD_BYTES: + self._cancel_pending_start( + "pending START frame exceeds maximum size" + ) + self.counters["dropped_malformed"] += 1 + continue + sent, failure = self._try_send_payload(payload) + if not sent: + self.counters["start_send_failures"] += 1 + self._cancel_pending_start(failure) + self.counters["frames_suppressed_inactive"] += 1 + continue + if merged: + self.counters["command_hand_merges"] += 1 + self.start_markers_remaining -= 1 + self._commit_pending_start_sent(now) + self._drain_responses() + self.write_status(force=True) + continue + if not self.teleop_active and transition != "stop": self.counters["frames_suppressed_inactive"] += 1 self._drain_responses() @@ -735,10 +1119,10 @@ class XteleSender: self.counters["frames_suppressed_inactive"] += 1 continue - if self.session is None and not self.connect(): + if self.session is None: self._abort_teleop( - "cannot start teleoperation because OmniSocket Hub is " - "unavailable; release Z+C before retrying" + "active OmniSocket session is unavailable; a new Z+C " + "cycle is required" ) self.counters["frames_suppressed_inactive"] += 1 continue @@ -753,44 +1137,58 @@ class XteleSender: else: session_state = "active" stop_reason = "" + stop_requested = transition == "stop" try: - payload, merged = self._build_payload( - data, - command, - session_id, - self.teleop_session_seq, - session_state, - stop_reason, - suppress_processed_right, - ) - except ValueError: - self.counters["dropped_malformed"] += 1 - continue - if merged: - self.counters["command_hand_merges"] += 1 - if len(payload) > MAX_PAYLOAD_BYTES: - self.counters["dropped_malformed"] += 1 - continue + try: + payload, merged = self._build_payload( + data, + command, + session_id, + self.teleop_session_seq, + session_state, + stop_reason, + suppress_processed_right, + ) + except ValueError: + self.counters["dropped_malformed"] += 1 + if stop_requested: + self.counters["teleop_stop_send_failures"] += 1 + self.last_error = "cannot encode operator STOP frame" + continue + if merged: + self.counters["command_hand_merges"] += 1 + if len(payload) > MAX_PAYLOAD_BYTES: + self.counters["dropped_malformed"] += 1 + if stop_requested: + self.counters["teleop_stop_send_failures"] += 1 + self.last_error = ( + "operator STOP frame exceeds maximum size" + ) + continue - sent = self._send_payload(payload) - if sent and session_state == "start": - self.start_markers_remaining -= 1 - if transition == "stop": - # STOP is the final xTELE business frame. The underlying - # Session is flushed for bounded delivery and then closed; - # the next physical START creates a fresh registration. - self.teleop_session_id = None - self.teleop_session_seq = 0 - self.start_markers_remaining = 0 - self._flush_session() - self.close_session() + sent = self._send_payload(payload) + if not sent and stop_requested: + self.counters["teleop_stop_send_failures"] += 1 + if sent and session_state == "start": + self.start_markers_remaining -= 1 + finally: + if stop_requested: + # STOP is final even when encode/send fails. The robot + # then falls back to its input timeout; EAI must never + # retain a registered, locally inactive Session. + self._finish_stop_session() self._drain_responses() self.write_status() finally: + self._invalidate_connect_attempt() source.close() if command_source is not None: command_source.close() context.term() + thread = self.connect_thread + if thread is not None: + thread.join(timeout=3.5) + self._poll_connect_result(time.monotonic()) self.close_session() self.write_status(force=True) return 0 diff --git a/tg3_omnisocket_transport/test_session_gate.py b/tg3_omnisocket_transport/test_session_gate.py index 6a2ba2c..b50cf32 100644 --- a/tg3_omnisocket_transport/test_session_gate.py +++ b/tg3_omnisocket_transport/test_session_gate.py @@ -7,8 +7,10 @@ import importlib.util import json from pathlib import Path import sys +from threading import Event from types import ModuleType import unittest +from unittest import mock fake_omnisocket = ModuleType("omnisocket") @@ -53,6 +55,71 @@ def frame(pressed: bool) -> dict[str, object]: class SessionGateTest(unittest.TestCase): + class FakeSession: + def __init__( + self, + responses: list[tuple[str, int, bytes]] | None = None, + ) -> None: + self.responses = list(responses or []) + self.sent: list[tuple[str, bytes]] = [] + self.closed = False + + @staticmethod + def stats() -> dict[str, int]: + return {"connected": 1, "registered": 1} + + @staticmethod + def kcp_stats() -> dict[str, int]: + return { + "snd_queue": 0, + "snd_buffer": 0, + "last_feedback_age_ms": 0, + } + + def send(self, *, to: str, data: bytes) -> None: + self.sent.append((to, data)) + + def recv(self, timeout_ms: int) -> tuple[str, int, bytes] | None: + self.assert_zero_timeout(timeout_ms) + if self.responses: + return self.responses.pop(0) + return None + + @staticmethod + def assert_zero_timeout(timeout_ms: int) -> None: + if timeout_ms != 0: + raise AssertionError(f"unexpected receive timeout: {timeout_ms}") + + def close(self) -> None: + self.closed = True + + class BlockingConnectSession(FakeSession): + def __init__( + self, + entered: Event, + release: Event, + failure: OSError | None = None, + ) -> None: + super().__init__() + self.entered = entered + self.release = release + self.failure = failure + self.connect_args: tuple[str, str, dict[str, object]] | None = None + + def connect( + self, + *, + server_addr: str, + peer_id: str, + **options: object, + ) -> None: + self.connect_args = (server_addr, peer_id, options) + self.entered.set() + if not self.release.wait(timeout=1.0): + raise RuntimeError("test did not release blocked connect") + if self.failure is not None: + raise self.failure + def setUp(self) -> None: self.sender = sender_module.XteleSender(args()) @@ -67,7 +134,43 @@ class SessionGateTest(unittest.TestCase): self.assertFalse(self.sender.require_combo_release) return finished_at - def test_boot_requires_stable_release_before_start(self) -> None: + def create_pending_start(self, started_at: float = 0.0) -> float: + released_at = self.stable_release(started_at) + self.assertIsNone( + self.sender._update_teleop_gate(released_at + 0.01, frame(True)) + ) + pending_at = released_at + 3.02 + self.assertEqual( + self.sender._update_teleop_gate(pending_at, frame(True)), + "start_pending", + ) + self.assertTrue(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + return pending_at + + def cancel_pending_from_latest_frame(self) -> None: + if ( + self.sender.start_pending + and self.sender.latest_combo_state is not True + ): + reason = ( + "released" + if self.sender.latest_combo_state is False + else "malformed" + ) + self.sender._cancel_pending_start(reason) + + def attach_session( + self, + responses: list[tuple[str, int, bytes]] | None = None, + ) -> FakeSession: + session = self.FakeSession(responses) + self.sender.session = session + self.sender.session_connected_at = 100.0 + self.sender.counters["connected"] = 1 + return session + + def test_boot_requires_stable_release_before_pending_start(self) -> None: self.assertIsNone(self.sender._update_teleop_gate(0.0, frame(True))) self.assertIsNone(self.sender._update_teleop_gate(4.0, frame(True))) released_at = self.stable_release(5.0) @@ -76,44 +179,574 @@ class SessionGateTest(unittest.TestCase): ) self.assertEqual( self.sender._update_teleop_gate(released_at + 3.02, frame(True)), - "start", + "start_pending", ) - self.assertTrue(self.sender.teleop_active) + self.assertTrue(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertEqual(self.sender.counters["teleop_start_requests"], 1) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + self.assertIsNotNone(self.sender.teleop_session_id) + self.assertTrue( + self.sender._pending_connect_ready(released_at + 3.02) + ) + + def test_boot_malformed_buttons_cannot_clear_release_requirement(self) -> None: + malformed = frame(False) + malformed["button"] = { + "left": [0, 0, 0], + "right": [0, 0, "0"], + } + + self.assertIsNone(self.sender._update_teleop_gate(0.0, frame(False))) + self.assertEqual(self.sender.combo_release_started_at, 0.0) + self.assertIsNone(self.sender._update_teleop_gate(0.49, malformed)) + self.assertTrue(self.sender.require_combo_release) + self.assertIsNone(self.sender.combo_release_started_at) + self.assertIsNone(self.sender.latest_combo_state) + + # The release timer must restart after malformed input; time before the + # malformed frame cannot be accumulated toward re-arming. + self.assertIsNone(self.sender._update_teleop_gate(0.50, frame(False))) + self.assertIsNone(self.sender._update_teleop_gate(0.99, frame(False))) + self.assertTrue(self.sender.require_combo_release) + self.assertIsNone(self.sender._update_teleop_gate(1.01, frame(False))) + self.assertFalse(self.sender.require_combo_release) + + def test_start_stop_parser_accepts_only_explicit_boolean_buttons(self) -> None: + self.assertTrue(self.sender._start_stop_pressed(frame(True))) + self.assertFalse(self.sender._start_stop_pressed(frame(False))) + for invalid in (1.0, "1", None, 2, -1): + with self.subTest(value=invalid): + malformed = frame(True) + malformed["button"] = { + "left": [0, 0, invalid], + "right": [0, 0, 1], + } + self.assertIsNone( + self.sender._start_stop_pressed(malformed) + ) def test_single_false_frame_cannot_rearm_stop(self) -> None: - released_at = self.stable_release(0.0) - self.sender._update_teleop_gate(released_at + 0.01, frame(True)) - self.assertEqual( - self.sender._update_teleop_gate(released_at + 3.02, frame(True)), - "start", - ) + pending_at = self.create_pending_start() + self.sender._commit_pending_start_sent(pending_at + 0.01) self.assertIsNone( - self.sender._update_teleop_gate(released_at + 3.03, frame(False)) + self.sender._update_teleop_gate(pending_at + 0.02, frame(False)) ) self.assertIsNone( - self.sender._update_teleop_gate(released_at + 3.04, frame(True)) + self.sender._update_teleop_gate(pending_at + 0.03, frame(True)) ) self.assertIsNone( - self.sender._update_teleop_gate(released_at + 7.00, frame(True)) + self.sender._update_teleop_gate(pending_at + 4.00, frame(True)) ) self.assertTrue(self.sender.teleop_active) self.assertTrue(self.sender.require_combo_release) def test_stable_release_allows_separate_stop_hold(self) -> None: - released_at = self.stable_release(0.0) - self.sender._update_teleop_gate(released_at + 0.01, frame(True)) - self.assertEqual( - self.sender._update_teleop_gate(released_at + 3.02, frame(True)), - "start", - ) - second_release = self.stable_release(released_at + 3.03) + pending_at = self.create_pending_start() + self.sender._commit_pending_start_sent(pending_at + 0.01) + second_release = self.stable_release(pending_at + 0.02) self.sender._update_teleop_gate(second_release + 0.01, frame(True)) self.assertEqual( self.sender._update_teleop_gate(second_release + 3.02, frame(True)), "stop", ) self.assertFalse(self.sender.teleop_active) + self.assertEqual(self.sender.counters["teleop_starts"], 1) + self.assertEqual(self.sender.counters["teleop_stops"], 1) + + def test_cancelled_start_requires_release_before_a_new_hold(self) -> None: + pending_at = self.create_pending_start() + first_session_id = self.sender.teleop_session_id + self.sender._cancel_pending_start("Hub offline") + + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertTrue(self.sender.require_combo_release) + self.assertEqual(self.sender.last_start_failure, "Hub offline") + self.assertIsNone(self.sender.teleop_session_id) + + # A consumed hold cannot silently become another pending request while + # the operator keeps Z+C pressed, regardless of elapsed time. + self.assertIsNone( + self.sender._update_teleop_gate( + pending_at + 10.0, frame(True) + ) + ) + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender._pending_connect_ready(pending_at + 10.0)) + + # A stable release followed by a new complete hold creates a distinct + # logical session which can make one new connection attempt. + released_at = self.stable_release(pending_at + 10.01) + self.assertIsNone( + self.sender._update_teleop_gate( + released_at + 0.01, frame(True) + ) + ) + self.assertEqual( + self.sender._update_teleop_gate( + released_at + 3.02, frame(True) + ), + "start_pending", + ) + self.assertTrue(self.sender.start_pending) + self.assertIsNotNone(self.sender.teleop_session_id) + self.assertNotEqual(self.sender.teleop_session_id, first_session_id) + self.assertEqual(self.sender.counters["teleop_start_requests"], 2) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + + def test_pending_start_release_cancels_and_closes_session(self) -> None: + pending_at = self.create_pending_start() + session = self.attach_session() + self.sender.start_wait_fresh_after_connect = True + + self.assertIsNone( + self.sender._update_teleop_gate(pending_at + 0.01, frame(False)) + ) + self.cancel_pending_from_latest_frame() + + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertIsNone(self.sender.session) + self.assertTrue(session.closed) + self.assertEqual(session.sent, []) + self.assertEqual(self.sender.counters["start_pending_cancels"], 1) + self.assertTrue(self.sender.require_combo_release) + self.assertFalse(self.sender._pending_connect_ready(pending_at + 10.0)) + + def test_pending_start_malformed_button_state_cancels(self) -> None: + pending_at = self.create_pending_start() + malformed = frame(True) + malformed["button"] = {"left": [0, 0], "right": [0, 0, 1]} + + self.assertIsNone( + self.sender._update_teleop_gate(pending_at + 0.01, malformed) + ) + self.cancel_pending_from_latest_frame() + + self.assertFalse(self.sender.start_pending) + self.assertIsNone(self.sender.latest_combo_state) + self.assertEqual(self.sender.last_start_failure, "malformed") + self.assertEqual(self.sender.counters["start_pending_cancels"], 1) + + def test_pending_freshness_waits_then_accepts_advancing_timestamp(self) -> None: + self.create_pending_start() + self.sender.start_wait_fresh_after_connect = True + first = frame(True) + first["timestamp"] = 100.0 + second = frame(True) + second["timestamp"] = 100.01 + + self.assertEqual( + self.sender._check_pending_start_freshness(first), + ("wait", ""), + ) + self.assertEqual(self.sender.start_fresh_barrier_timestamp, 100.0) + self.assertTrue(self.sender.start_wait_fresh_after_connect) + self.assertFalse(self.sender.teleop_active) + + self.assertEqual( + self.sender._check_pending_start_freshness(second), + ("ready", ""), + ) + self.assertIsNone(self.sender.start_fresh_barrier_timestamp) + self.assertFalse(self.sender.start_wait_fresh_after_connect) + self.assertTrue(self.sender.start_pending) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + + def test_pending_freshness_missing_timestamp_cancels(self) -> None: + self.create_pending_start() + session = self.attach_session() + self.sender.start_wait_fresh_after_connect = True + + outcome, reason = self.sender._check_pending_start_freshness( + frame(True) + ) + self.assertEqual(outcome, "cancel") + self.assertIn("fresh xTELE timestamp", reason) + self.sender._cancel_pending_start(reason) + + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertTrue(session.closed) + self.assertEqual(session.sent, []) + + def test_pending_freshness_repeated_timestamp_waits_for_later_advance( + self, + ) -> None: + self.create_pending_start() + self.sender.start_wait_fresh_after_connect = True + first = frame(True) + first["timestamp"] = 100.0 + repeated = frame(True) + repeated["timestamp"] = 100.0 + advanced = frame(True) + advanced["timestamp"] = 100.01 + + self.assertEqual( + self.sender._check_pending_start_freshness(first), + ("wait", ""), + ) + self.assertEqual( + self.sender._check_pending_start_freshness(repeated), + ("wait", ""), + ) + self.assertEqual(self.sender.start_fresh_barrier_timestamp, 100.0) + self.assertTrue(self.sender.start_wait_fresh_after_connect) + self.assertTrue(self.sender.start_pending) + self.assertEqual( + self.sender._check_pending_start_freshness(advanced), + ("ready", ""), + ) + self.assertIsNone(self.sender.start_fresh_barrier_timestamp) + self.assertFalse(self.sender.start_wait_fresh_after_connect) + self.assertTrue(self.sender.start_pending) + + def test_pending_freshness_backwards_timestamp_cancels(self) -> None: + self.create_pending_start() + session = self.attach_session() + self.sender.start_wait_fresh_after_connect = True + first = frame(True) + first["timestamp"] = 100.0 + backwards = frame(True) + backwards["timestamp"] = 99.99 + + self.assertEqual( + self.sender._check_pending_start_freshness(first), + ("wait", ""), + ) + outcome, reason = self.sender._check_pending_start_freshness(backwards) + self.assertEqual(outcome, "cancel") + self.assertIn("moved backwards", reason) + self.sender._cancel_pending_start(reason) + + self.assertFalse(self.sender.start_pending) + self.assertTrue(session.closed) + self.assertEqual(session.sent, []) + + def test_background_connect_success_is_adopted_as_pending(self) -> None: + self.create_pending_start() + entered = Event() + release = Event() + session = self.BlockingConnectSession(entered, release) + + with ( + mock.patch.object(sender_module, "Session", return_value=session), + mock.patch.object(self.sender, "write_status"), + ): + self.assertTrue(self.sender._start_connect_attempt()) + self.assertTrue(entered.wait(timeout=1.0)) + self.assertTrue(self.sender._connect_inflight()) + thread = self.sender.connect_thread + self.assertIsNotNone(thread) + release.set() + thread.join(timeout=1.0) # type: ignore[union-attr] + self.assertFalse(thread.is_alive()) # type: ignore[union-attr] + self.assertEqual( + self.sender._poll_connect_result(50.0), "connected" + ) + + self.assertIs(self.sender.session, session) + self.assertFalse(session.closed) + self.assertTrue(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertTrue(self.sender.start_wait_fresh_after_connect) + self.assertIsNone(self.sender.start_fresh_barrier_timestamp) + self.assertEqual(self.sender.counters["connected"], 1) + self.assertEqual(self.sender.counters["reconnects"], 1) + self.assertEqual( + session.connect_args, + (self.sender.args.server, self.sender.args.peer_id, {}), + ) + self.sender.close_session() + + def test_release_during_connect_discards_and_closes_late_result(self) -> None: + pending_at = self.create_pending_start() + entered = Event() + release = Event() + session = self.BlockingConnectSession(entered, release) + + with mock.patch.object(sender_module, "Session", return_value=session): + self.assertTrue(self.sender._start_connect_attempt()) + self.assertTrue(entered.wait(timeout=1.0)) + thread = self.sender.connect_thread + self.assertIsNotNone(thread) + + self.assertIsNone( + self.sender._update_teleop_gate( + pending_at + 0.01, frame(False) + ) + ) + self.cancel_pending_from_latest_frame() + self.assertFalse(self.sender.start_pending) + + release.set() + thread.join(timeout=1.0) # type: ignore[union-attr] + self.assertFalse(thread.is_alive()) # type: ignore[union-attr] + self.assertEqual( + self.sender._poll_connect_result(pending_at + 0.02), + "discarded", + ) + + self.assertTrue(session.closed) + self.assertIsNone(self.sender.session) + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertEqual(self.sender.counters["connected"], 0) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + self.assertEqual(self.sender.counters["start_pending_cancels"], 1) + + def test_background_connect_failure_consumes_hold_without_retry(self) -> None: + pending_at = self.create_pending_start() + first_session_id = self.sender.teleop_session_id + entered = Event() + release = Event() + session = self.BlockingConnectSession( + entered, + release, + OSError("connection refused"), + ) + + with ( + mock.patch.object(sender_module, "Session", return_value=session), + mock.patch.object(self.sender, "write_status"), + ): + # run() owns this counter; mirror its one increment around the + # lower-level worker call used by this unit test. + self.sender.counters["start_connect_attempts"] += 1 + self.assertTrue(self.sender._start_connect_attempt()) + self.assertTrue(entered.wait(timeout=1.0)) + thread = self.sender.connect_thread + self.assertIsNotNone(thread) + release.set() + thread.join(timeout=1.0) # type: ignore[union-attr] + self.assertFalse(thread.is_alive()) # type: ignore[union-attr] + self.assertEqual(self.sender._poll_connect_result(20.0), "failed") + + self.assertTrue(session.closed) + self.assertEqual(session.sent, []) + self.assertIsNone(self.sender.session) + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertTrue(self.sender.require_combo_release) + self.assertEqual(self.sender.counters["start_connect_failures"], 1) + self.assertEqual(self.sender.counters["start_connect_attempts"], 1) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + self.assertEqual(self.sender.counters["frames_sent"], 0) + self.assertIn("connection refused", self.sender.last_start_failure) + self.assertFalse(self.sender._pending_connect_ready(100.0)) + + # Keeping the failed Z+C hold pressed cannot create another request or + # connection attempt after the Hub later becomes available. + self.assertIsNone( + self.sender._update_teleop_gate(25.0, frame(True)) + ) + self.assertFalse(self.sender.start_pending) + self.assertEqual(self.sender.counters["teleop_start_requests"], 1) + self.assertEqual(self.sender.counters["start_connect_attempts"], 1) + + released_at = self.stable_release(25.01) + self.assertIsNone( + self.sender._update_teleop_gate( + released_at + 0.01, frame(True) + ) + ) + self.assertEqual( + self.sender._update_teleop_gate( + released_at + 3.02, frame(True) + ), + "start_pending", + ) + self.assertTrue(self.sender.start_pending) + self.assertNotEqual(self.sender.teleop_session_id, first_session_id) + self.assertEqual(self.sender.counters["teleop_start_requests"], 2) + self.assertTrue(self.sender._pending_connect_ready(released_at + 3.02)) + + def test_session_constructor_failure_consumes_hold_without_retry(self) -> None: + self.create_pending_start() + constructor_called = Event() + + def broken_session_constructor() -> object: + constructor_called.set() + raise RuntimeError("Session constructor failed") + + with ( + mock.patch.object( + sender_module, + "Session", + side_effect=broken_session_constructor, + ), + mock.patch.object(self.sender, "write_status"), + ): + self.assertTrue(self.sender._start_connect_attempt()) + self.assertTrue(constructor_called.wait(timeout=1.0)) + thread = self.sender.connect_thread + self.assertIsNotNone(thread) + thread.join(timeout=1.0) # type: ignore[union-attr] + self.assertFalse(thread.is_alive()) # type: ignore[union-attr] + self.assertEqual(self.sender._poll_connect_result(30.0), "failed") + + self.assertIsNone(self.sender.session) + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertTrue(self.sender.require_combo_release) + self.assertEqual(self.sender.counters["start_connect_failures"], 1) + self.assertIn("Session constructor failed", self.sender.last_start_failure) + self.assertEqual(self.sender.counters["teleop_starts"], 0) + self.assertFalse(self.sender._pending_connect_ready(300.0)) + + def test_successful_start_commit_activates_once(self) -> None: + pending_at = self.create_pending_start() + self.sender._commit_pending_start_sent(pending_at + 0.01) + + self.assertFalse(self.sender.start_pending) + self.assertTrue(self.sender.teleop_active) + self.assertEqual(self.sender.counters["teleop_starts"], 1) + self.assertEqual(self.sender.counters["teleop_start_requests"], 1) + + def test_unknown_target_is_an_ordinary_remote_error_and_aborts(self) -> None: + pending_at = self.create_pending_start() + self.sender._commit_pending_start_sent(pending_at + 0.01) + session = self.attach_session( + [("hub", fake_omnisocket.MSG_TYPE_ERROR, b"unknown target: robot")] + ) + + self.sender._drain_responses() + + self.assertTrue(session.closed) + self.assertIsNone(self.sender.session) + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertTrue(self.sender.require_combo_release) + self.assertEqual(self.sender.counters["teleop_aborts"], 1) + self.assertEqual(self.sender.counters["teleop_starts"], 1) + self.assertIn("unknown target: robot", self.sender.last_error) + + def test_other_remote_error_also_aborts(self) -> None: + pending_at = self.create_pending_start() + self.sender._commit_pending_start_sent(pending_at + 0.01) + session = self.attach_session( + [("hub", fake_omnisocket.MSG_TYPE_ERROR, b"route unavailable")] + ) + + self.sender._drain_responses() + + self.assertTrue(session.closed) + self.assertFalse(self.sender.start_pending) + self.assertFalse(self.sender.teleop_active) + self.assertEqual(self.sender.counters["teleop_aborts"], 1) + self.assertIn("route unavailable", self.sender.last_error) + + def test_source_clock_allows_adjacent_repeat_under_timeout(self) -> None: + first = frame(False) + first["timestamp"] = 2507 + repeated = frame(False) + repeated["timestamp"] = 2507 + + self.assertEqual( + self.sender._check_source_clock(10.0, first), + (True, ""), + ) + self.assertEqual( + self.sender._check_source_clock(10.24, repeated), + (True, ""), + ) + self.assertEqual(self.sender.last_source_timestamp, 2507.0) + self.assertEqual(self.sender.last_source_timestamp_advanced_at, 10.0) + + def test_source_clock_rejects_persistent_freeze_over_timeout(self) -> None: + first = frame(False) + first["timestamp"] = 2517 + repeated = frame(False) + repeated["timestamp"] = 2517 + + self.assertEqual( + self.sender._check_source_clock(10.0, first), + (True, ""), + ) + self.assertEqual( + self.sender._check_source_clock(10.24, repeated), + (True, ""), + ) + clock_ok, reason = self.sender._check_source_clock(10.251, repeated) + self.assertFalse(clock_ok) + self.assertIn("stopped advancing", reason) + self.assertEqual(self.sender.last_source_timestamp, 2517.0) + + def test_source_clock_rejects_backwards_timestamp_while_active(self) -> None: + first = frame(False) + first["timestamp"] = 3000 + backwards = frame(False) + backwards["timestamp"] = 2999 + self.assertEqual( + self.sender._check_source_clock(10.0, first), + (True, ""), + ) + self.sender.teleop_active = True + + clock_ok, reason = self.sender._check_source_clock(10.01, backwards) + + self.assertFalse(clock_ok) + self.assertIn("moved backwards", reason) + self.assertEqual(self.sender.last_source_timestamp, 3000.0) + self.assertEqual(self.sender.counters["source_timestamp_resets"], 0) + + def test_source_clock_idle_backwards_resets_and_requires_release(self) -> None: + first = frame(False) + first["timestamp"] = 4000 + restarted = frame(True) + restarted["timestamp"] = 10 + self.assertEqual( + self.sender._check_source_clock(10.0, first), + (True, ""), + ) + self.sender.require_combo_release = False + self.sender.combo_started_at = 10.01 + self.sender.combo_release_started_at = 10.02 + self.sender.latest_combo_state = True + + self.assertEqual( + self.sender._check_source_clock(10.03, restarted), + (True, ""), + ) + + self.assertEqual(self.sender.last_source_timestamp, 10.0) + self.assertEqual(self.sender.last_source_timestamp_advanced_at, 10.03) + self.assertTrue(self.sender.require_combo_release) + self.assertIsNone(self.sender.combo_started_at) + self.assertIsNone(self.sender.combo_release_started_at) + self.assertIsNone(self.sender.latest_combo_state) + self.assertEqual(self.sender.counters["source_timestamp_resets"], 1) + + def test_source_clock_rejects_missing_timestamp(self) -> None: + clock_ok, reason = self.sender._check_source_clock(10.0, frame(False)) + + self.assertFalse(clock_ok) + self.assertIn("missing or malformed", reason) + self.assertIsNone(self.sender.last_source_timestamp) + + def test_finish_stop_session_clears_state_and_closes_session(self) -> None: + session = self.attach_session() + self.sender.teleop_active = True + self.sender.teleop_session_id = "active-session" + self.sender.teleop_session_seq = 42 + self.sender.start_markers_remaining = 7 + + self.sender._finish_stop_session() + + self.assertFalse(self.sender.teleop_active) + self.assertIsNone(self.sender.teleop_session_id) + self.assertEqual(self.sender.teleop_session_seq, 0) + self.assertEqual(self.sender.start_markers_remaining, 0) + self.assertIsNone(self.sender.session) + self.assertTrue(session.closed) + self.assertEqual(session.sent, []) + self.assertEqual(self.sender.counters["connected"], 0) def test_untrusted_transport_metadata_is_overwritten(self) -> None: data = frame(False)