From ba0c148dabdcc1367de14adff6d786e558af5386 Mon Sep 17 00:00:00 2001 From: HappyDog0713 Date: Tue, 15 Sep 2026 07:58:08 +0000 Subject: [PATCH 1/4] refactor(vla): retire legacy RoboTwin inference endpoint Remove the LingBot-specific RoboTwin MessagePack server, validator, dependency file, and dead action scheduler compatibility alias. Keep direct inference, native HTTP, generic VLA session WebSocket, and MuJoCo entrypoints as the supported integration paths. Update VLA documentation and the MuJoCo example to use the generic service defaults. Verification: scripts/run_ci_tests.sh --skip-install (lint, format, import, CPU unit, and server tests); local MuJoCo smoke; focused LingBot VLA and VLA session tests. --- docs/en/vla.md | 29 +- examples/lingbot_vla_v2/README.md | 164 +----- .../lingbot_vla_v2/lingbot_vla_v2_mujoco.py | 2 +- .../lingbot_vla_v2_robotwin_server.py | 452 --------------- .../lingbot_vla_v2/requirements-robotwin.txt | 2 - .../lingbot_vla_v2/action_scheduler.py | 5 - .../lingbot_vla_v2/test_action_scheduler.py | 131 ----- .../lingbot_vla_v2/test_robotwin_server.py | 273 --------- .../test_lingbot_vla_v2_robotwin_ws.py | 182 ------ .../validate_lingbot_vla_v2_robotwin_ws.py | 537 ------------------ 10 files changed, 20 insertions(+), 1757 deletions(-) delete mode 100644 examples/lingbot_vla_v2/lingbot_vla_v2_robotwin_server.py delete mode 100644 examples/lingbot_vla_v2/requirements-robotwin.txt delete mode 100644 telefuser/pipelines/lingbot_vla_v2/action_scheduler.py delete mode 100644 tests/unit/pipelines/lingbot_vla_v2/test_action_scheduler.py delete mode 100644 tests/unit/pipelines/lingbot_vla_v2/test_robotwin_server.py delete mode 100644 tests/unit/validation/test_lingbot_vla_v2_robotwin_ws.py delete mode 100644 tools/validation/validate_lingbot_vla_v2_robotwin_ws.py diff --git a/docs/en/vla.md b/docs/en/vla.md index 4801f03..f718476 100644 --- a/docs/en/vla.md +++ b/docs/en/vla.md @@ -40,7 +40,7 @@ The public modules are: - `telefuser.service.vla_replica`: worker-local OPEN, PREDICT, RESET, and CLOSE dispatch for pipeline replicas. - `telefuser.client.AsyncVLAClient`: remote session client with concurrent request correlation. -## LingBot-VLA v2 Compatibility +## LingBot-VLA v2 Integration `LingBotVlaV2VLAPolicy` wraps `LingBotVlaV2Pipeline`; the pipeline's existing tensor input and return types are not changed. It labels the normalized canonical `[T,55]` result as a `ModelActionChunk`. @@ -49,16 +49,9 @@ changed. It labels the normalized canonical `[T,55]` result as a `ModelActionChu chunk to an absolute-position `[H,14]` `RobotActionChunk` in the declared dual-arm joint order. Its model and robot action spaces are exported as `LINGBOT_VLA_V2_ACTION_SPACE` and `ROBOTWIN_ACTION_SPACE`. -The generic VLA WebSocket is the primary online integration path. The existing LingBot RoboTwin WebSocket endpoint -is a compatibility adapter for unmodified upstream `WebsocketClientPolicy` clients and calls the same policy, -embodiment, session, and runtime path internally. Its URL, MessagePack request fields, metadata frame, response fields, -reset behavior, and latest-wins -scheduler behavior remain compatible. The old model-specific scheduler import is retained as an alias to the common -runtime scheduler. - -The standalone endpoint owns one resident pipeline, so each connection session is inherently pinned to that policy -instance. The shared HTTP structured-task route remains unchanged and continues to return the existing canonical -JSON result. +The generic VLA WebSocket is the online integration path for both RoboTwin and MuJoCo clients. Each connection opens +an explicit model and embodiment pair, while the shared HTTP structured-task route remains unchanged and continues to +return the existing canonical JSON result. ## Session Lifecycle @@ -102,9 +95,7 @@ robot action and observation spaces; semantic incompatibility is rejected before The protocol returns machine-readable error codes for malformed messages, unsupported versions, unknown components or sessions, contract mismatches, superseded work, out-of-order or expired observations, request timeout, unavailable sessions, and replica failure. Tensor payloads are dense, typed, shape-checked Base64 data inside bounded JSON -messages. This is the stable interoperability format; the LingBot-RoboTwin compatibility endpoint continues to use -its existing MessagePack format. New model and simulator integrations must target the generic protocol rather than -add behavior to the model-specific compatibility endpoint. +messages. This is the stable interoperability format for all model and simulator integrations. `request_ttl_ms` is measured with server monotonic time and bounds inference delivery. Observation age is checked separately using `observation_timestamp_ns` and `observation_clock_now_ns`, which must come from the same clock domain. @@ -160,14 +151,10 @@ route. The LingBot example supplies a standalone server that starts `PipelinePoo This keeps existing HTTP routing and every non-VLA pipeline unchanged. A real simulator still owns control timing, actuator feedback, and verification that each returned action was actually applied. -For LingBot-VLA v2, the standalone generic WebSocket is the primary continuous-control entrypoint. The direct Python +For LingBot-VLA v2, the standalone generic WebSocket is the continuous-control entrypoint. The direct Python entrypoint remains the reference/offline baseline, and the native HTTP structured service remains for existing -TeleFuser callers. The RoboTwin MessagePack server is a legacy compatibility entrypoint only; it can be removed after -all upstream clients migrate to `/v1/vla/session`. - -The direct inference CLI and structured HTTP task remain separate because they provide reference and batch workflows, -not simulator session transports. The legacy RoboTwin MessagePack endpoint should be removed only after the RTX client -passes generic-protocol action delivery, reset, timeout, reconnect, and latest-wins parity checks. +TeleFuser callers. The direct inference CLI and structured HTTP task remain separate because they provide reference +and batch workflows, not simulator session transports. No new dependencies, environment variables, shared model configuration fields, CLI options, HTTP schemas, or existing service routes are introduced by this semantic and transport layer. diff --git a/examples/lingbot_vla_v2/README.md b/examples/lingbot_vla_v2/README.md index bea2d69..3ef0a75 100644 --- a/examples/lingbot_vla_v2/README.md +++ b/examples/lingbot_vla_v2/README.md @@ -25,7 +25,6 @@ The parity reference uses [Robbyant/lingbot-vla-v2](https://github.com/Robbyant/ | CUDA Graph | Supported | Dynamic eager prefix with an opt-in fixed-shape action-denoising graph | | Quantization | Partial | Profile-specific release status; see Configuration and Performance | | Native server API | Supported | Asynchronous structured task API and `TFClient` | -| RoboTwin policy protocol | Compatibility | Legacy MessagePack endpoint for unmodified upstream clients | | Request replicas | Supported | One complete policy copy per GPU | | Single-policy FSDP, TP, or PP | Unsupported | The integration does not split one policy across GPUs | | RoboTwin action mapping | Supported | Unnormalizes canonical output to absolute-position `50 x 14` chunks | @@ -192,15 +191,14 @@ Compare a deterministic quantized capture with the corresponding TeleFuser BF16 ## Serving -The generic VLA session server is the only recommended online path for new simulator integrations. The four current -entrypoints share one `LingBotVlaV2Pipeline`; they are access modes, not separate model implementations: +The three inference entrypoints share one `LingBotVlaV2Pipeline`; they are access modes, not separate model +implementations: | Entry point | Role | Recommendation | | --- | --- | --- | | `lingbot_vla_v2_inference.py` | Direct Python reference and offline baseline | Keep for regression/debugging | | `lingbot_vla_v2_native_service.py` via `telefuser serve` | Native HTTP structured requests | Keep for TeleFuser compatibility | | `lingbot_vla_v2_vla_server.py` | Stateful generic VLA WebSocket | **Primary simulator path** | -| `lingbot_vla_v2_robotwin_server.py` | Upstream RoboTwin MessagePack compatibility | Legacy; remove after client migration | For a production or simulation deployment, start only the generic WebSocket server. The other entries remain for baseline comparison and backward compatibility and do not change the model or session implementation. @@ -268,156 +266,16 @@ For a minimal P0 check, use one fake adapter test for report status and one asyn long-running soak, cross-machine clock comparison, and full simulator episode are not required to validate this runtime API. -### Legacy RoboTwin Protocol Compatibility +### RoboTwin via Generic VLA Session -This compatibility server implements the persistent MessagePack WebSocket protocol used by the upstream -`WebsocketClientPolicy`. Use it only when the upstream client cannot yet consume the generic VLA session protocol. -New transport, scheduling, and simulator integrations belong on the generic VLA path. The compatibility endpoint is -isolated from `telefuser serve`: no TeleFuser API routes, service schemas, or other model integrations are changed. +RoboTwin clients use the same `/v1/vla/session` endpoint as MuJoCo. The XPolicyLab-side `TeleFuser_VLA` policy maps +RoboTwin observations to the negotiated `robotwin` embodiment contract and returns validated absolute-position action +chunks. No model-specific server or MessagePack dependency is required in TeleFuser. -Install the protocol dependency in the TeleFuser inference environment: - -```bash -.venv-vla/bin/python -m pip install -r examples/lingbot_vla_v2/requirements-robotwin.txt -``` - -Start one resident policy process: - -```bash -.venv-vla/bin/python examples/lingbot_vla_v2/lingbot_vla_v2_robotwin_server.py \ - --model-root "$TF_MODEL_ZOO_PATH/lingbot/lingbot-vla-v2-6b" \ - --qwen3vl-root "$TF_MODEL_ZOO_PATH/Qwen3-VL-4B-Instruct" \ - --device cuda:0 --host 0.0.0.0 --port 9330 --use-length 50 -``` - -On NVIDIA H100, this dedicated entrypoint disables cuDNN SDPA before model warmup because the current -PyTorch/cuDNN combination cannot build a valid vision-attention execution plan. Flash, memory-efficient, and math -SDPA remain enabled. The override is process-local and is not applied to other TeleFuser pipelines. - -The server exposes `GET /healthz` and the policy WebSocket at `/`. On connection it sends a MessagePack metadata -frame, including the explicit 16 MiB request limit, then accepts multiple binary MessagePack requests on the same -connection. This matches the upstream client contract: - -```python -from deploy.websocket_client_policy import WebsocketClientPolicy - -policy = WebsocketClientPolicy(host="127.0.0.1", port=9330) -policy.reset("robotwin") -result = policy.infer( - { - "observation.images.cam_high": camera_high, - "observation.images.cam_left_wrist": camera_left_wrist, - "observation.images.cam_right_wrist": camera_right_wrist, - "observation.state": state, - "task": instruction, - } -) -actions = result["action"] # float32 [50, 14] when --use-length=50 -``` - -The initial metadata frame describes protocol version `1.0`, `absolute_qpos` action semantics, `float32` dtype, -horizon, dimension, and the exact dual-arm joint order. Inference requests may include an integer `seed` plus -`request_id` and `episode_id`; the response echoes them and reports decode, lock-wait, pipeline, action-mapping, and -adapter timings. Existing clients may omit all three request fields. - -The endpoint also advertises an additive, latest-wins action scheduler. A client that overlaps simulation and -inference should send a monotonically increasing `sequence_id` within each `episode_id`, plus a positive -`request_ttl_ms`. The server has one GPU worker, retains at most one pending request per connection/episode, and -accepts new observations while inference is running. A newer observation replaces queued work; because an in-flight -CUDA call cannot be cancelled, its result is discarded after completion when it has become stale. Successful -responses use `scheduler_status="completed"`. Responses with `superseded`, `expired`, `stale_sequence`, or -`overloaded` contain `action=None` and a structured `error`; clients must never execute those responses. - -`request_ttl_ms` starts when the H100 server receives the request. Do not compare monotonic timestamps between the -H100 and RTX machines. The RTX client should separately enforce its round-trip deadline and hold the current joint -positions when no fresh action is available. - -Validate this direct endpoint before a simulator is available. This sends reset and repeated inference requests to -the resident model, validates the returned `[H, 14]` action contract, and optionally verifies exact fixed-seed replay: - -```bash -.venv-vla/bin/python -m tools.validation.validate_lingbot_vla_v2_robotwin_ws \ - --host 127.0.0.1 --port 9330 \ - --image examples/data/lingbot_world_fast/image.jpg \ - --max-image-edge 640 \ - --task "pick up the object" --seed 7 --requests 10 \ - --output work_dirs/robotwin_ws_validation/smoke.json -``` - -The validator preserves aspect ratio and downsizes only images whose longest edge exceeds `--max-image-edge`, then -checks the encoded MessagePack request against the limit advertised by the server before sending it. This keeps the -large repository sample representative of normal RoboTwin camera payloads. Add `--require-exact-replay` only when -validating a runtime profile that promises bitwise determinism; BF16 H100 inference is validated with numerical -tolerances rather than identical action hashes. - -Exercise overlapping submissions and stale-action rejection without a simulator: - -```bash -.venv-vla/bin/python -m tools.validation.validate_lingbot_vla_v2_robotwin_ws \ - --host 127.0.0.1 --port 9330 \ - --image examples/data/lingbot_world_fast/image.jpg \ - --task "pick up the object" --seed 7 --requests 3 \ - --request-ttl-ms 5000 --overlap-requests \ - --output work_dirs/robotwin_ws_validation/overlap.json -``` - -This mode sends all observations before receiving responses, requires the newest request to return an action, and -requires at least one older request to be reported as `superseded`. - -Each request runs the existing pipeline, converts normalized canonical `50 x 55` output through the bundled RoboTwin -profile, and returns absolute-position actions in raw RoboTwin order. `--use-length` may truncate the returned chunk; -start with 50 for upstream-equivalent open-loop execution. The adapter accepts episode reset messages but deliberately -rejects runtime checkpoint switching. - -Internally this compatibility endpoint uses the shared `VLAPolicy`, `EmbodimentAdapter`, `VLASessionManager`, and -`ChunkExecutor` contracts. The wire protocol and the existing `LingBotVlaV2Pipeline` API remain unchanged. The -declared control rate is intentionally unresolved until the remote RoboTwin loop supplies its actual frequency. - -Keep this endpoint until the RTX client has passed end-to-end action delivery, reset, timeout, reconnect, and -latest-wins parity checks through the generic protocol. After that migration, the compatibility module can be removed -without changing the model pipeline or the generic VLA service. - -For split-machine deployment, run the model endpoint and the repository-owned XPolicyLab proxy on the H100 inference -host. The proxy does not load a second model; it translates XPolicyLab observations to the direct TeleFuser protocol: - -```bash -cd /data/RoboTwin -bash XPolicyLab/policy/TeleFuser_LingBot_VLA/setup_eval_policy_server.sh \ - RoboTwin lift_pot remote_base arx_x5 joint 0 0 \ - /data/RoboTwin/.venv 19000 0.0.0.0 \ - 127.0.0.1 9330 -``` - -On the remote RTX/Vulkan workstation, use the standard RoboTwin evaluation client and point it at the proxy. No -TeleFuser files or model weights are required on that workstation: - -```bash -cd /data/RoboTwin -bash scripts/eval_policy.sh \ - --bench_name RoboTwin \ - --task_name lift_pot \ - --env_cfg_type arx_x5 \ - --policy_name TeleFuser_LingBot_VLA \ - --host INFERENCE_HOST --port 19000 --protocol ws \ - --eval_batch false --root_dir /data/RoboTwin --device_id 0 \ - --additional_info ckpt_name=remote_base,action_type=joint \ - --seed 0 --task_config demo_clean --test_num 1 -``` - -The current XPolicyLab proxy calls `infer()` synchronously, so it remains compatible but does not yet overlap action -execution with inference. Full overlap requires an incremental RTX-side change: execute chunk N while submitting a -newer observation for chunk N+1, keep only the newest completed chunk in an atomic action buffer, and apply the same -sequence/deadline checks before execution. That simulator-side change is outside this repository and is not required -for the no-simulation server validation above. - -Keep ports `9330` and `19000` on a trusted private network or an SSH/VPN tunnel. These WebSocket endpoints do not -provide authentication or transport encryption. The direct validator covers preprocessing, inference, mapping, and -the inner WebSocket contract; only the RTX smoke episode can additionally establish XPolicyLab translation and one -real SAPIEN simulation step. - -The base checkpoint remains marked `unverified_official_6b_base`. This endpoint establishes preprocessing, inference, -action mapping, transport, and simulator execution continuity; it does not establish RoboTwin task success without -an embodiment-validated checkpoint. +Run the generic server on the inference host, then point the XPolicyLab adapter at +`ws://INFERENCE_HOST:8000/v1/vla/session` with `model_id=lingbot-vla-v2` and `embodiment_id=robotwin`. Keep the generic +WebSocket and the XPolicyLab policy server on a trusted private network or an SSH/VPN tunnel because neither endpoint +provides authentication. ## Local MuJoCo smoke simulation @@ -436,7 +294,7 @@ bash examples/lingbot_vla_v2/setup_mujoco_local.sh # Complete local simulator -> generic VLA WebSocket -> simulator loop .venv/bin/python examples/lingbot_vla_v2/lingbot_vla_v2_mujoco.py \ - --mode websocket --server-url ws://127.0.0.1:18080/v1/vla/session \ + --mode websocket --server-url ws://127.0.0.1:8000/v1/vla/session \ --chunks 2 --execute-horizon 8 \ --output-dir work_dirs/lingbot_vla_v2/mujoco_websocket ``` diff --git a/examples/lingbot_vla_v2/lingbot_vla_v2_mujoco.py b/examples/lingbot_vla_v2/lingbot_vla_v2_mujoco.py index 36cd24d..b17cc69 100644 --- a/examples/lingbot_vla_v2/lingbot_vla_v2_mujoco.py +++ b/examples/lingbot_vla_v2/lingbot_vla_v2_mujoco.py @@ -236,7 +236,7 @@ async def _run_websocket_loop( @click.command() @click.option("--mode", type=click.Choice(("local", "websocket")), default="local", show_default=True) @click.option("--urdf", "urdf_path", type=click.Path(path_type=Path, exists=True, dir_okay=False), default=DEFAULT_URDF) -@click.option("--server-url", default="ws://127.0.0.1:18080/v1/vla/session", show_default=True) +@click.option("--server-url", default="ws://127.0.0.1:8000/v1/vla/session", show_default=True) @click.option("--instruction", default="pick up the red block", show_default=True) @click.option("--seed", type=int, default=7, show_default=True) @click.option("--chunks", type=click.IntRange(min=1), default=2, show_default=True) diff --git a/examples/lingbot_vla_v2/lingbot_vla_v2_robotwin_server.py b/examples/lingbot_vla_v2/lingbot_vla_v2_robotwin_server.py deleted file mode 100644 index ec28f8e..0000000 --- a/examples/lingbot_vla_v2/lingbot_vla_v2_robotwin_server.py +++ /dev/null @@ -1,452 +0,0 @@ -"""Legacy compatibility server for the upstream RoboTwin policy protocol. - -New simulator integrations should use the generic VLA WebSocket protocol in -``lingbot_vla_v2_vla_server.py``. This endpoint remains only for unmodified -upstream ``WebsocketClientPolicy`` clients. -""" - -from __future__ import annotations - -import asyncio -import contextlib -import operator -import threading -import time -from contextlib import asynccontextmanager -from typing import Any, Mapping, Protocol - -import click -import msgpack -import numpy as np -import torch -import uvicorn -from fastapi import FastAPI, WebSocket, WebSocketDisconnect - -from telefuser.pipelines.lingbot_vla_v2 import ( - ROBOTWIN_ACTION_ORDER, - ROBOTWIN_CAMERA_KEYS, - RobotWinProfile, - create_lingbot_vla_v2_session_manager, -) -from telefuser.pipelines.lingbot_vla_v2.runtime import ( - LINGBOT_VLA_V2_QUANTIZATION_CHOICES, - configure_lingbot_vla_v2_h100_sdpa, - get_lingbot_vla_v2_pipeline, -) -from telefuser.utils.logging import logger -from telefuser.vla import RobotObservation, RobotState -from telefuser.vla.runtime import ActionChunkScheduler - -ROBOTWIN_PROTOCOL_VERSION = "1.0" -ROBOTWIN_ACTION_TYPE = "absolute_qpos" -ROBOTWIN_ACTION_DTYPE = "float32" -ROBOTWIN_MAX_REQUEST_BYTES = 16 * 1024 * 1024 -_TRACE_ID_FIELDS = ("request_id", "episode_id") -_VLA_SESSION_ID_FIELD = "_telefuser_vla_session_id" - - -class _Pipeline(Protocol): - config: Any - - def __call__(self, observation: Any, seed: int | None = None) -> Any: ... - - def close(self) -> None: ... - - -def _pack_numpy(value: Any) -> Any: - """Encode NumPy values using the upstream msgpack_numpy wire format.""" - if isinstance(value, (np.ndarray, np.generic)) and value.dtype.kind in ("V", "O", "c"): - raise ValueError(f"unsupported NumPy dtype: {value.dtype}") - if isinstance(value, np.ndarray): - return { - b"__ndarray__": True, - b"data": value.tobytes(), - b"dtype": value.dtype.str, - b"shape": value.shape, - } - if isinstance(value, np.generic): - return { - b"__npgeneric__": True, - b"data": value.item(), - b"dtype": value.dtype.str, - } - raise TypeError(f"cannot encode value of type {type(value)!r}") - - -def _unpack_numpy(value: dict[Any, Any]) -> Any: - """Decode NumPy values produced by the upstream msgpack_numpy helper.""" - if b"__ndarray__" in value: - return np.ndarray( - buffer=value[b"data"], - dtype=np.dtype(value[b"dtype"]), - shape=value[b"shape"], - ) - if b"__npgeneric__" in value: - return np.dtype(value[b"dtype"]).type(value[b"data"]) - return value - - -def pack_message(payload: Mapping[str, Any]) -> bytes: - """Pack one RoboTwin policy protocol message.""" - return msgpack.packb(dict(payload), default=_pack_numpy) - - -def unpack_message(payload: bytes) -> dict[str, Any]: - """Unpack and validate one RoboTwin policy protocol message.""" - decoded = msgpack.unpackb(payload, object_hook=_unpack_numpy, raw=False) - if not isinstance(decoded, dict): - raise ValueError("RoboTwin request must be a MessagePack object") - return decoded - - -def _optional_seed(request: Mapping[str, Any]) -> int | None: - value = request.get("seed") - if value is None: - return None - if isinstance(value, bool): - raise ValueError("seed must be an integer") - try: - return operator.index(value) - except TypeError as error: - raise ValueError("seed must be an integer") from error - - -def _trace_fields(request: Mapping[str, Any]) -> dict[str, str | int]: - fields: dict[str, str | int] = {} - for name in _TRACE_ID_FIELDS: - value = request.get(name) - if value is None: - continue - if isinstance(value, bool) or not isinstance(value, str | int) or isinstance(value, str) and not value: - raise ValueError(f"{name} must be a non-empty string or integer") - fields[name] = value - return fields - - -def _optional_nonnegative_int(request: Mapping[str, Any], field: str) -> int | None: - value = request.get(field) - if value is None: - return None - if isinstance(value, bool): - raise ValueError(f"{field} must be a non-negative integer") - try: - parsed = operator.index(value) - except TypeError as error: - raise ValueError(f"{field} must be a non-negative integer") from error - if parsed < 0: - raise ValueError(f"{field} must be a non-negative integer") - return parsed - - -class RobotWinPolicyAdapter: - """Translate upstream RoboTwin observations to the TeleFuser VLA SDK.""" - - def __init__( - self, - pipeline: _Pipeline, - *, - profile: RobotWinProfile | None = None, - use_length: int = 50, - ) -> None: - if not 1 <= use_length <= 50: - raise ValueError(f"use_length must be in [1, 50], got {use_length}") - self.pipeline = pipeline - self.profile = profile or pipeline.config.robot_profile - self.use_length = use_length - self._lock = threading.Lock() - self._next_sequence: dict[str, int] = {} - self._sessions = create_lingbot_vla_v2_session_manager(pipeline, profile=self.profile) - - @property - def metadata(self) -> dict[str, Any]: - """Describe the action contract sent when a client connects.""" - return { - "protocol_version": ROBOTWIN_PROTOCOL_VERSION, - "robot_profile": self.profile.name, - "action_type": ROBOTWIN_ACTION_TYPE, - "action_horizon": self.use_length, - "action_dim": self.profile.raw_state_dim, - "action_dtype": ROBOTWIN_ACTION_DTYPE, - "action_order": list(ROBOTWIN_ACTION_ORDER), - "max_request_bytes": ROBOTWIN_MAX_REQUEST_BYTES, - "policy_verified": False, - "verification_status": "unverified_official_6b_base", - } - - def infer(self, request: Mapping[str, Any]) -> dict[str, Any]: - """Return one absolute-position RoboTwin action chunk.""" - trace_fields = _trace_fields(request) - if request.get("reset", False): - return {**self._reset(request), **trace_fields} - - missing = [key for key in (*ROBOTWIN_CAMERA_KEYS, "observation.state", "task") if key not in request] - if missing: - raise ValueError(f"RoboTwin observation is missing fields: {missing}") - - seed = _optional_seed(request) - episode_id = str(trace_fields.get("episode_id", "default")) - session_id_value = request.get(_VLA_SESSION_ID_FIELD, f"legacy:{episode_id}") - if not isinstance(session_id_value, str) or not session_id_value: - raise ValueError("internal VLA session ID must be a non-empty string") - sequence_id = _optional_nonnegative_int(request, "sequence_id") - observation_timestamp_ns = _optional_nonnegative_int(request, "observation_timestamp_ns") - if observation_timestamp_ns is None: - observation_timestamp_ns = time.time_ns() - robot_observation = RobotObservation( - state=RobotState( - values=torch.tensor(request["observation.state"], dtype=torch.float32, device="cpu"), - dimension_names=ROBOTWIN_ACTION_ORDER, - timestamp_ns=observation_timestamp_ns, - ), - images={key: request[key] for key in ROBOTWIN_CAMERA_KEYS}, - ) - adapter_started_at = time.monotonic() - with self._lock: - lock_wait_ms = (time.monotonic() - adapter_started_at) * 1000.0 - if sequence_id is None: - sequence_id = self._next_sequence.get(session_id_value, 0) - self._next_sequence[session_id_value] = sequence_id + 1 - if session_id_value not in self._sessions.session_ids(): - self._sessions.open( - session_id_value, - model_id="lingbot-vla-v2", - embodiment_id=self.profile.embodiment_id, - episode_id=episode_id, - execute_horizon=self.use_length, - ) - vla_timings: dict[str, float] = {} - action_chunk = self._sessions.get(session_id_value).predict( - robot_observation, - request["task"], - sequence_id, - seed=seed, - timings=vla_timings, - ) - pipeline_ms = vla_timings["policy_ms"] - action_mapping_ms = vla_timings["decode_actions_ms"] + vla_timings["prepare_actions_ms"] - if action_chunk.valid_length < self.use_length: - raise RuntimeError( - f"policy returned horizon {action_chunk.valid_length}, shorter than use_length={self.use_length}" - ) - actions = np.ascontiguousarray(action_chunk.actions[: self.use_length].numpy(), dtype=np.float32) - expected_shape = (self.use_length, self.profile.raw_state_dim) - if actions.shape != expected_shape: - raise RuntimeError(f"mapped actions must have shape {expected_shape}, got {actions.shape}") - if not np.isfinite(actions).all(): - raise RuntimeError("mapped actions must contain only finite values") - response: dict[str, Any] = { - "action": actions, - "policy_verified": action_chunk.metadata["policy_verified"], - "verification_status": action_chunk.metadata["verification_status"], - "server_timing": { - "lock_wait_ms": lock_wait_ms, - "pipeline_ms": pipeline_ms, - "action_mapping_ms": action_mapping_ms, - "adapter_total_ms": (time.monotonic() - adapter_started_at) * 1000.0, - }, - **trace_fields, - } - if seed is not None: - response["seed"] = seed - return response - - def _reset(self, request: Mapping[str, Any]) -> dict[str, Any]: - robot_name = request.get("robo_name", self.profile.name) - if robot_name != self.profile.name: - raise ValueError(f"unsupported robot profile: {robot_name!r}") - if request.get("path_to_pi_model") not in (None, ""): - raise ValueError("runtime checkpoint switching is not supported") - episode_id = str(request.get("episode_id", "default")) - session_id = request.get(_VLA_SESSION_ID_FIELD, f"legacy:{episode_id}") - if isinstance(session_id, str) and session_id in self._sessions.session_ids(): - self._sessions.reset(session_id, episode_id) - self._next_sequence[session_id] = 0 - return {"action": None} - - def release_session(self, session_id: str) -> None: - """Release semantic session state after a WebSocket disconnect.""" - with self._lock: - if session_id in self._sessions.session_ids(): - self._sessions.close(session_id) - self._next_sequence.pop(session_id, None) - - def close(self) -> None: - """Release resources owned by the resident policy.""" - for session_id in self._sessions.session_ids(): - self._sessions.close(session_id) - self._next_sequence.clear() - self.pipeline.close() - - -def create_robotwin_app( - adapter: RobotWinPolicyAdapter, - *, - max_pending_sessions: int = 32, -) -> FastAPI: - """Create a standalone app compatible with upstream WebsocketClientPolicy.""" - scheduler = ActionChunkScheduler(adapter.infer, max_pending_sessions=max_pending_sessions) - - @asynccontextmanager - async def lifespan(_app: FastAPI): - await scheduler.start() - try: - yield - finally: - await scheduler.close() - - app = FastAPI(title="LingBot-VLA v2 RoboTwin Policy", lifespan=lifespan) - - @app.get("/healthz") - async def healthz() -> dict[str, str]: - return {"status": "ok"} - - @app.websocket("/") - async def policy_socket(websocket: WebSocket) -> None: - await websocket.accept() - await websocket.send_bytes(pack_message({**adapter.metadata, **scheduler.metadata})) - connection_key = str(id(websocket)) - session_keys: set[str] = set() - delivery_tasks: set[asyncio.Task[None]] = set() - send_lock = asyncio.Lock() - previous_total_ms: float | None = None - - async def deliver( - response_future: asyncio.Future[dict[str, Any]], - *, - decode_ms: float, - round_started_at: float, - ) -> None: - nonlocal previous_total_ms - try: - response = dict(await response_future) - server_timing = dict(response.get("server_timing", {})) - server_timing["decode_ms"] = decode_ms - response["server_timing"] = server_timing - async with send_lock: - if previous_total_ms is not None: - server_timing["prev_total_ms"] = previous_total_ms - await websocket.send_bytes(pack_message(response)) - previous_total_ms = (time.monotonic() - round_started_at) * 1000.0 - except asyncio.CancelledError: - raise - except WebSocketDisconnect: - return - except Exception as error: - logger.exception("LingBot-VLA v2 RoboTwin request failed") - async with send_lock: - with contextlib.suppress(WebSocketDisconnect, RuntimeError): - await websocket.send_text(f"{type(error).__name__}: {error}") - await websocket.close(code=1011) - - try: - while True: - message = await websocket.receive() - if message["type"] == "websocket.disconnect": - break - payload = message.get("bytes") - if payload is None: - raise ValueError("RoboTwin requests must use binary MessagePack frames") - - round_started_at = time.monotonic() - decode_started_at = time.monotonic() - request = unpack_message(payload) - decode_ms = (time.monotonic() - decode_started_at) * 1000.0 - trace_fields = _trace_fields(request) - episode_id = trace_fields.get("episode_id", "default") - session_key = f"{connection_key}:{episode_id!r}" - request[_VLA_SESSION_ID_FIELD] = session_key - if request.get("reset", False): - for previous_session_key in session_keys - {session_key}: - scheduler.release_session(previous_session_key) - await asyncio.to_thread(adapter.release_session, previous_session_key) - scheduler.release_session(session_key) - session_keys.intersection_update({session_key}) - session_keys.add(session_key) - response_future = scheduler.submit(request, session_key=session_key) - task = asyncio.create_task( - deliver( - response_future, - decode_ms=decode_ms, - round_started_at=round_started_at, - ) - ) - delivery_tasks.add(task) - task.add_done_callback(delivery_tasks.discard) - except WebSocketDisconnect: - return - except Exception as error: - logger.exception("LingBot-VLA v2 RoboTwin request failed") - with contextlib.suppress(WebSocketDisconnect, RuntimeError): - await websocket.send_text(f"{type(error).__name__}: {error}") - await websocket.close(code=1011) - finally: - for session_key in session_keys: - scheduler.release_session(session_key) - for task in delivery_tasks: - task.cancel() - if delivery_tasks: - await asyncio.gather(*delivery_tasks, return_exceptions=True) - if session_keys: - await asyncio.gather( - *(asyncio.to_thread(adapter.release_session, session_key) for session_key in session_keys) - ) - - return app - - -@click.command() -@click.option("--model-root", required=True, type=click.Path(exists=True, file_okay=False)) -@click.option("--qwen3vl-root", required=True, type=click.Path(exists=True, file_okay=False)) -@click.option("--host", default="0.0.0.0", show_default=True) -@click.option("--port", default=9330, show_default=True, type=click.IntRange(1, 65535)) -@click.option("--device", default="cuda:0", show_default=True) -@click.option("--use-length", default=50, show_default=True, type=click.IntRange(1, 50)) -@click.option( - "--max-pending-sessions", - default=32, - show_default=True, - type=click.IntRange(1), - help="Bound the number of sessions waiting behind the GPU worker", -) -@click.option("--cuda-graph", is_flag=True, help="Enable fixed-shape CUDA Graph inference") -@click.option( - "--quantization", - type=click.Choice(LINGBOT_VLA_V2_QUANTIZATION_CHOICES), - default=None, -) -def main( - model_root: str, - qwen3vl_root: str, - host: str, - port: int, - device: str, - use_length: int, - max_pending_sessions: int, - cuda_graph: bool, - quantization: str | None, -) -> None: - """Start the legacy-compatible policy endpoint for an upstream RoboTwin client.""" - configure_lingbot_vla_v2_h100_sdpa(device) - pipeline = get_lingbot_vla_v2_pipeline( - model_root, - qwen3vl_root, - device=device, - warmup=True, - quantization=quantization, - cuda_graph=cuda_graph, - ) - adapter = RobotWinPolicyAdapter(pipeline, use_length=use_length) - try: - uvicorn.run( - create_robotwin_app(adapter, max_pending_sessions=max_pending_sessions), - host=host, - port=port, - workers=1, - ws_max_size=ROBOTWIN_MAX_REQUEST_BYTES, - ) - finally: - adapter.close() - - -if __name__ == "__main__": - main() diff --git a/examples/lingbot_vla_v2/requirements-robotwin.txt b/examples/lingbot_vla_v2/requirements-robotwin.txt deleted file mode 100644 index 7beb7e9..0000000 --- a/examples/lingbot_vla_v2/requirements-robotwin.txt +++ /dev/null @@ -1,2 +0,0 @@ -msgpack>=1.0,<2.0 -websockets>=12.0,<16.0 diff --git a/telefuser/pipelines/lingbot_vla_v2/action_scheduler.py b/telefuser/pipelines/lingbot_vla_v2/action_scheduler.py deleted file mode 100644 index 5427a74..0000000 --- a/telefuser/pipelines/lingbot_vla_v2/action_scheduler.py +++ /dev/null @@ -1,5 +0,0 @@ -"""Compatibility import for the generic VLA action scheduler.""" - -from telefuser.vla.runtime.scheduler import ActionChunkScheduler - -__all__ = ["ActionChunkScheduler"] diff --git a/tests/unit/pipelines/lingbot_vla_v2/test_action_scheduler.py b/tests/unit/pipelines/lingbot_vla_v2/test_action_scheduler.py deleted file mode 100644 index fcd5c6f..0000000 --- a/tests/unit/pipelines/lingbot_vla_v2/test_action_scheduler.py +++ /dev/null @@ -1,131 +0,0 @@ -from __future__ import annotations - -import asyncio -import threading -import time -from typing import Any, Mapping - -import pytest - -from telefuser.pipelines.lingbot_vla_v2.action_scheduler import ActionChunkScheduler -from telefuser.vla.runtime import ActionChunkScheduler as CommonActionChunkScheduler - - -def test_lingbot_scheduler_import_is_a_compatibility_alias() -> None: - assert ActionChunkScheduler is CommonActionChunkScheduler - - -def test_scheduler_discards_inflight_and_pending_work_when_newer_observation_arrives() -> None: - started = threading.Event() - release = threading.Event() - - def infer(request: Mapping[str, Any]) -> dict[str, Any]: - if request["sequence_id"] == 0: - started.set() - assert release.wait(timeout=2) - return {"action": request["sequence_id"]} - - async def scenario() -> None: - scheduler = ActionChunkScheduler(infer) - await scheduler.start() - first = scheduler.submit({"request_id": "first", "sequence_id": 0}, session_key="session") - assert await asyncio.to_thread(started.wait, 2) - second = scheduler.submit({"request_id": "second", "sequence_id": 1}, session_key="session") - third = scheduler.submit({"request_id": "third", "sequence_id": 2}, session_key="session") - - assert (await second)["scheduler_status"] == "superseded" - release.set() - assert (await first)["scheduler_status"] == "superseded" - completed = await third - assert completed["action"] == 2 - assert completed["sequence_id"] == 2 - assert completed["scheduler_status"] == "completed" - assert completed["server_timing"]["queue_wait_ms"] >= 0 - await scheduler.close() - - asyncio.run(scenario()) - - -def test_scheduler_expires_result_after_non_cancellable_inference() -> None: - def infer(_request: Mapping[str, Any]) -> dict[str, Any]: - time.sleep(0.02) - return {"action": "too-late"} - - async def scenario() -> None: - scheduler = ActionChunkScheduler(infer) - await scheduler.start() - response = await scheduler.submit( - {"request_id": "expired", "request_ttl_ms": 1}, - session_key="session", - ) - assert response["action"] is None - assert response["scheduler_status"] == "expired" - assert response["error"]["code"] == "expired" - await scheduler.close() - - asyncio.run(scenario()) - - -def test_scheduler_rejects_stale_sequence_without_running_inference() -> None: - calls: list[int] = [] - - def infer(request: Mapping[str, Any]) -> dict[str, Any]: - calls.append(request["sequence_id"]) - return {"action": request["sequence_id"]} - - async def scenario() -> None: - scheduler = ActionChunkScheduler(infer) - await scheduler.start() - assert (await scheduler.submit({"sequence_id": 4}, session_key="session"))["action"] == 4 - stale = await scheduler.submit({"sequence_id": 4}, session_key="session") - assert stale["scheduler_status"] == "stale_sequence" - assert calls == [4] - await scheduler.close() - - asyncio.run(scenario()) - - -def test_scheduler_bounds_pending_sessions() -> None: - started = threading.Event() - release = threading.Event() - - def infer(request: Mapping[str, Any]) -> dict[str, Any]: - if request["request_id"] == "active": - started.set() - assert release.wait(timeout=2) - return {"action": request["request_id"]} - - async def scenario() -> None: - scheduler = ActionChunkScheduler(infer, max_pending_sessions=1) - await scheduler.start() - active = scheduler.submit({"request_id": "active"}, session_key="active") - assert await asyncio.to_thread(started.wait, 2) - pending = scheduler.submit({"request_id": "pending"}, session_key="pending") - overloaded = await scheduler.submit({"request_id": "overloaded"}, session_key="overloaded") - assert overloaded["scheduler_status"] == "overloaded" - release.set() - assert (await active)["scheduler_status"] == "completed" - assert (await pending)["scheduler_status"] == "completed" - await scheduler.close() - - asyncio.run(scenario()) - - -@pytest.mark.parametrize( - ("payload", "match"), - [ - ({"sequence_id": -1}, "sequence_id"), - ({"sequence_id": True}, "sequence_id"), - ({"request_ttl_ms": 0}, "request_ttl_ms"), - ({"request_ttl_ms": float("inf")}, "request_ttl_ms"), - ], -) -def test_scheduler_rejects_invalid_controls(payload: dict[str, Any], match: str) -> None: - async def scenario() -> None: - scheduler = ActionChunkScheduler(lambda _request: {"action": 1}) - await scheduler.start() - with pytest.raises(ValueError, match=match): - scheduler.submit(payload, session_key="session") - await scheduler.close() - - asyncio.run(scenario()) diff --git a/tests/unit/pipelines/lingbot_vla_v2/test_robotwin_server.py b/tests/unit/pipelines/lingbot_vla_v2/test_robotwin_server.py deleted file mode 100644 index 3ace520..0000000 --- a/tests/unit/pipelines/lingbot_vla_v2/test_robotwin_server.py +++ /dev/null @@ -1,273 +0,0 @@ -from __future__ import annotations - -import threading -from types import SimpleNamespace -from typing import Any - -import numpy as np -import pytest -import torch -from click.testing import CliRunner -from fastapi.testclient import TestClient - -from examples.lingbot_vla_v2 import lingbot_vla_v2_robotwin_server as server -from telefuser.pipelines.lingbot_vla_v2 import ( - ROBOTWIN_CAMERA_KEYS, - LingBotVlaV2CanonicalActionChunk, - RobotWinProfile, -) - - -def _profile() -> RobotWinProfile: - return RobotWinProfile( - { - "observation.state.arm.position": {"q01": [0.0] * 12, "q99": [2.0] * 12}, - "observation.state.effector.position": {"q01": [-1.0] * 2, "q99": [1.0] * 2}, - "action.arm.position": {"q01": [0.0] * 12, "q99": [2.0] * 12}, - "action.effector.position": {"q01": [-1.0] * 2, "q99": [1.0] * 2}, - } - ) - - -class _Pipeline: - def __init__(self, profile: RobotWinProfile) -> None: - self.config = SimpleNamespace(robot_profile=profile) - self.observations: list[Any] = [] - self.seeds: list[int | None] = [] - self.closed = False - - def __call__(self, observation, seed: int | None = None) -> LingBotVlaV2CanonicalActionChunk: - self.observations.append(observation) - self.seeds.append(seed) - return LingBotVlaV2CanonicalActionChunk( - canonical_normalized_actions=torch.zeros(50, 55), - horizon=50, - action_dim=55, - ) - - def close(self) -> None: - self.closed = True - - -def _observation() -> dict[str, Any]: - images = {key: np.zeros((8, 8, 3), dtype=np.uint8) for key in ROBOTWIN_CAMERA_KEYS} - return { - **images, - "observation.state": np.zeros(14, dtype=np.float32), - "task": "pick up the block", - } - - -def test_message_codec_matches_upstream_numpy_contract() -> None: - payload = { - "image": np.arange(24, dtype=np.uint8).reshape(2, 4, 3), - "state": np.float32(1.5), - } - - decoded = server.unpack_message(server.pack_message(payload)) - - assert np.array_equal(decoded["image"], payload["image"]) - assert decoded["image"].dtype == np.uint8 - assert decoded["state"] == np.float32(1.5) - - -def test_message_codec_rejects_unsafe_object_arrays() -> None: - with pytest.raises(ValueError, match="unsupported NumPy dtype"): - server.pack_message({"value": np.asarray([object()], dtype=object)}) - - -def test_adapter_maps_observation_and_returns_robotwin_chunk() -> None: - profile = _profile() - pipeline = _Pipeline(profile) - adapter = server.RobotWinPolicyAdapter(pipeline, use_length=8) - - request = {**_observation(), "seed": 7, "request_id": "request-1", "episode_id": "episode-1"} - result = adapter.infer(request) - - assert result["action"].shape == (8, 14) - assert result["action"].dtype == np.float32 - assert np.allclose(result["action"][:, [6, 13]], 0.0, atol=1e-6) - assert np.allclose(np.delete(result["action"], [6, 13], axis=1), 1.0000005) - assert result["policy_verified"] is False - assert result["seed"] == 7 - assert result["request_id"] == "request-1" - assert result["episode_id"] == "episode-1" - assert result["server_timing"]["lock_wait_ms"] >= 0 - assert result["server_timing"]["pipeline_ms"] >= 0 - assert result["server_timing"]["action_mapping_ms"] >= 0 - assert result["server_timing"]["adapter_total_ms"] >= 0 - assert len(pipeline.observations) == 1 - assert pipeline.seeds == [7] - observation = pipeline.observations[0] - assert observation.task == "pick up the block" - assert set(observation.images) == set(ROBOTWIN_CAMERA_KEYS) - - -def test_adapter_reset_does_not_run_or_reload_policy() -> None: - pipeline = _Pipeline(_profile()) - adapter = server.RobotWinPolicyAdapter(pipeline) - - assert adapter.infer({"reset": True, "robo_name": "robotwin"}) == {"action": None} - assert pipeline.observations == [] - with pytest.raises(ValueError, match="runtime checkpoint switching"): - adapter.infer({"reset": True, "path_to_pi_model": "/different/checkpoint"}) - - -@pytest.mark.parametrize("field", ["request_id", "episode_id"]) -def test_adapter_rejects_invalid_trace_fields(field: str) -> None: - adapter = server.RobotWinPolicyAdapter(_Pipeline(_profile())) - - with pytest.raises(ValueError, match=field): - adapter.infer({**_observation(), field: ""}) - - -@pytest.mark.parametrize("seed", [True, 1.5, "7"]) -def test_adapter_rejects_non_integer_seed(seed: Any) -> None: - adapter = server.RobotWinPolicyAdapter(_Pipeline(_profile())) - - with pytest.raises(ValueError, match="seed must be an integer"): - adapter.infer({**_observation(), "seed": seed}) - - -def test_adapter_rejects_missing_observation_fields() -> None: - adapter = server.RobotWinPolicyAdapter(_Pipeline(_profile())) - - with pytest.raises(ValueError, match="missing fields"): - adapter.infer({"task": "pick up the block"}) - - -def test_websocket_is_persistent_and_uses_upstream_response_fields() -> None: - pipeline = _Pipeline(_profile()) - adapter = server.RobotWinPolicyAdapter(pipeline, use_length=3) - - with TestClient(server.create_robotwin_app(adapter)) as client: - assert client.get("/healthz").json() == {"status": "ok"} - with client.websocket_connect("/") as websocket: - metadata = server.unpack_message(websocket.receive_bytes()) - assert metadata["robot_profile"] == "robotwin" - assert metadata["protocol_version"] == "1.0" - assert metadata["action_type"] == "absolute_qpos" - assert metadata["action_horizon"] == 3 - assert metadata["action_dim"] == 14 - assert metadata["action_dtype"] == "float32" - assert metadata["action_order"] == list(server.ROBOTWIN_ACTION_ORDER) - assert metadata["max_request_bytes"] == server.ROBOTWIN_MAX_REQUEST_BYTES - assert metadata["scheduling"]["mode"] == "latest_wins" - assert metadata["scheduling"]["max_pending_per_session"] == 1 - assert metadata["scheduling"]["inflight_cancellation"] is False - - websocket.send_bytes( - server.pack_message( - { - **_observation(), - "seed": 11, - "request_id": "request-1", - "episode_id": "episode-1", - "sequence_id": 0, - "request_ttl_ms": 5_000, - } - ) - ) - response = server.unpack_message(websocket.receive_bytes()) - assert response["action"].shape == (3, 14) - assert response["seed"] == 11 - assert response["request_id"] == "request-1" - assert response["episode_id"] == "episode-1" - assert response["sequence_id"] == 0 - assert response["scheduler_status"] == "completed" - assert response["server_timing"]["decode_ms"] >= 0 - assert response["server_timing"]["infer_ms"] >= 0 - assert response["server_timing"]["queue_wait_ms"] >= 0 - assert response["server_timing"]["scheduler_total_ms"] >= 0 - assert response["server_timing"]["pipeline_ms"] >= 0 - assert response["server_timing"]["action_mapping_ms"] >= 0 - - websocket.send_bytes( - server.pack_message( - { - "reset": True, - "robo_name": "robotwin", - "request_id": "reset-1", - "episode_id": "episode-1", - } - ) - ) - reset_response = server.unpack_message(websocket.receive_bytes()) - assert reset_response["action"] is None - assert reset_response["request_id"] == "reset-1" - assert reset_response["episode_id"] == "episode-1" - assert reset_response["server_timing"]["prev_total_ms"] >= 0 - - websocket.send_bytes( - server.pack_message( - { - **_observation(), - "request_id": "after-reset", - "episode_id": "episode-1", - "sequence_id": 0, - } - ) - ) - after_reset = server.unpack_message(websocket.receive_bytes()) - assert after_reset["scheduler_status"] == "completed" - assert after_reset["sequence_id"] == 0 - assert after_reset["action"].shape == (3, 14) - - -def test_websocket_accepts_overlapping_chunks_and_discards_superseded_actions() -> None: - first_started = threading.Event() - release_first = threading.Event() - call_count = 0 - - class BlockingPipeline(_Pipeline): - def __call__(self, observation, seed: int | None = None): - nonlocal call_count - call_count += 1 - if call_count == 1: - first_started.set() - assert release_first.wait(timeout=2) - return super().__call__(observation, seed=seed) - - pipeline = BlockingPipeline(_profile()) - adapter = server.RobotWinPolicyAdapter(pipeline, use_length=3) - - with TestClient(server.create_robotwin_app(adapter)) as client: - with client.websocket_connect("/") as websocket: - server.unpack_message(websocket.receive_bytes()) - for sequence_id in range(3): - websocket.send_bytes( - server.pack_message( - { - **_observation(), - "request_id": f"request-{sequence_id}", - "episode_id": "episode", - "sequence_id": sequence_id, - } - ) - ) - if sequence_id == 0: - assert first_started.wait(timeout=2) - release_first.set() - responses = { - response["request_id"]: response - for response in (server.unpack_message(websocket.receive_bytes()) for _ in range(3)) - } - - assert responses["request-0"]["scheduler_status"] == "superseded" - assert responses["request-0"]["action"] is None - assert responses["request-1"]["scheduler_status"] == "superseded" - assert responses["request-1"]["action"] is None - assert responses["request-2"]["scheduler_status"] == "completed" - assert responses["request-2"]["action"].shape == (3, 14) - assert call_count == 2 - - -def test_cli_exposes_isolated_robotwin_server_options() -> None: - result = CliRunner().invoke(server.main, ["--help"]) - - assert result.exit_code == 0 - assert "--model-root" in result.output - assert "--qwen3vl-root" in result.output - assert "--use-length" in result.output - assert "--max-pending-sessions" in result.output - assert "--cuda-graph" in result.output diff --git a/tests/unit/validation/test_lingbot_vla_v2_robotwin_ws.py b/tests/unit/validation/test_lingbot_vla_v2_robotwin_ws.py deleted file mode 100644 index 8014901..0000000 --- a/tests/unit/validation/test_lingbot_vla_v2_robotwin_ws.py +++ /dev/null @@ -1,182 +0,0 @@ -from __future__ import annotations - -import argparse - -import numpy as np -import pytest -from PIL import Image - -from tools.validation import validate_lingbot_vla_v2_robotwin_ws as validator - - -def _metadata() -> dict: - return { - "protocol_version": validator.PROTOCOL_VERSION, - "robot_profile": "robotwin", - "action_type": validator.ACTION_TYPE, - "action_horizon": 3, - "action_dim": len(validator.ACTION_ORDER), - "action_dtype": validator.ACTION_DTYPE, - "action_order": list(validator.ACTION_ORDER), - "policy_verified": False, - "verification_status": "unverified_official_6b_base", - "max_request_bytes": 16 * 1024 * 1024, - "scheduling": { - "mode": "latest_wins", - "max_pending_per_session": 1, - "max_pending_sessions": 32, - "sequence_field": "sequence_id", - "ttl_field": "request_ttl_ms", - "inflight_cancellation": False, - }, - } - - -def _response() -> dict: - return { - "action": np.arange(42, dtype=np.float32).reshape(3, 14), - "seed": 7, - "request_id": "request-1", - "episode_id": "episode-1", - "sequence_id": 3, - "scheduler_status": "completed", - "policy_verified": False, - "verification_status": "unverified_official_6b_base", - "server_timing": { - "decode_ms": 0.1, - "infer_ms": 2.0, - "queue_wait_ms": 0.1, - "scheduler_total_ms": 2.1, - "lock_wait_ms": 0.05, - "pipeline_ms": 1.5, - "action_mapping_ms": 0.2, - "adapter_total_ms": 1.8, - }, - } - - -def test_validate_metadata_accepts_contract_and_additive_fields() -> None: - metadata = {**_metadata(), "future_field": "ignored"} - - summary = validator.validate_metadata(metadata) - - assert summary["action_shape"] == [3, 14] - assert summary["action_type"] == "absolute_qpos" - - -@pytest.mark.parametrize( - ("field", "value"), - [ - ("protocol_version", "2.0"), - ("action_type", "delta_qpos"), - ("action_dtype", "float64"), - ("action_order", ["unknown"] * 14), - ("action_horizon", 0), - ], -) -def test_validate_metadata_rejects_contract_mismatch(field: str, value) -> None: - metadata = {**_metadata(), field: value} - - with pytest.raises(validator.ValidationFailure, match=field): - validator.validate_metadata(metadata) - - -def test_validate_action_response_returns_stable_float32_digest() -> None: - summary = validator.validate_action_response( - _response(), - expected_horizon=3, - request_id="request-1", - episode_id="episode-1", - sequence_id=3, - seed=7, - ) - - assert summary["shape"] == [3, 14] - assert summary["dtype"] == "float32" - assert len(summary["sha256_float32_le"]) == 64 - assert summary["server_timing_ms"]["pipeline_ms"] == 1.5 - - -@pytest.mark.parametrize( - "actions", - [ - np.zeros((2, 14), dtype=np.float32), - np.zeros((3, 14), dtype=np.float64), - np.full((3, 14), np.nan, dtype=np.float32), - ], -) -def test_validate_action_response_rejects_invalid_action(actions: np.ndarray) -> None: - response = {**_response(), "action": actions} - - with pytest.raises(validator.ValidationFailure, match="action"): - validator.validate_action_response( - response, - expected_horizon=3, - request_id="request-1", - episode_id="episode-1", - sequence_id=3, - seed=7, - ) - - -def test_validate_reset_response_checks_trace_and_timings() -> None: - summary = validator.validate_reset_response( - { - "action": None, - "request_id": "reset-1", - "episode_id": "episode-1", - "server_timing": {"decode_ms": 0.1, "infer_ms": 0.2}, - }, - request_id="reset-1", - episode_id="episode-1", - ) - - assert summary["server_timing_ms"] == {"decode_ms": 0.1, "infer_ms": 0.2} - - -def test_validate_discarded_response_rejects_action_execution() -> None: - summary = validator.validate_discarded_response( - { - "action": None, - "request_id": "request-1", - "episode_id": "episode-1", - "sequence_id": 3, - "scheduler_status": "superseded", - "error": {"code": "superseded", "message": "newer observation received"}, - "server_timing": {"decode_ms": 0.1}, - }, - request_id="request-1", - episode_id="episode-1", - sequence_id=3, - ) - - assert summary["scheduler_status"] == "superseded" - - -def test_require_exact_replay_rejects_divergent_digests() -> None: - records = [ - {"action": {"sha256_float32_le": "a"}}, - {"action": {"sha256_float32_le": "b"}}, - ] - - with pytest.raises(validator.ValidationFailure, match="replay diverged"): - validator.require_exact_replay(records) - - -def test_parse_state_json_rejects_non_finite_state() -> None: - with pytest.raises(argparse.ArgumentTypeError, match="finite numbers"): - validator.parse_state_json("[0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 1e999]") - - -def test_load_validation_image_bounds_longest_edge_without_upscaling(tmp_path) -> None: - large_path = tmp_path / "large.png" - small_path = tmp_path / "small.png" - Image.new("RGB", (800, 400)).save(large_path) - Image.new("RGB", (32, 16)).save(small_path) - - large, source_shape = validator.load_validation_image(large_path, max_image_edge=640) - small, _ = validator.load_validation_image(small_path, max_image_edge=640) - - assert source_shape == [400, 800, 3] - assert large.shape == (320, 640, 3) - assert small.shape == (16, 32, 3) diff --git a/tools/validation/validate_lingbot_vla_v2_robotwin_ws.py b/tools/validation/validate_lingbot_vla_v2_robotwin_ws.py deleted file mode 100644 index 738874a..0000000 --- a/tools/validation/validate_lingbot_vla_v2_robotwin_ws.py +++ /dev/null @@ -1,537 +0,0 @@ -"""Validate a running LingBot-VLA v2 RoboTwin WebSocket endpoint without simulation.""" - -from __future__ import annotations - -import argparse -import hashlib -import json -import math -import statistics -import time -from pathlib import Path -from typing import Any, Mapping - -import numpy as np -from PIL import Image - -PROTOCOL_VERSION = "1.0" -ACTION_TYPE = "absolute_qpos" -ACTION_DTYPE = "float32" -DEFAULT_MAX_IMAGE_EDGE = 640 -ACTION_ORDER = ( - "left_arm_joint_0", - "left_arm_joint_1", - "left_arm_joint_2", - "left_arm_joint_3", - "left_arm_joint_4", - "left_arm_joint_5", - "left_gripper", - "right_arm_joint_0", - "right_arm_joint_1", - "right_arm_joint_2", - "right_arm_joint_3", - "right_arm_joint_4", - "right_arm_joint_5", - "right_gripper", -) -CAMERA_KEYS = ( - "observation.images.cam_high", - "observation.images.cam_left_wrist", - "observation.images.cam_right_wrist", -) -_ACTION_TIMING_FIELDS = ( - "decode_ms", - "infer_ms", - "queue_wait_ms", - "scheduler_total_ms", - "lock_wait_ms", - "pipeline_ms", - "action_mapping_ms", - "adapter_total_ms", -) - - -class ValidationFailure(RuntimeError): - """Raised when the endpoint violates the RoboTwin action contract.""" - - -def parse_state_json(value: str) -> list[float]: - """Parse a finite 14-dimensional RoboTwin state.""" - try: - raw = json.loads(value) - except json.JSONDecodeError as error: - raise argparse.ArgumentTypeError("state must be valid JSON") from error - if not isinstance(raw, list) or len(raw) != 14: - raise argparse.ArgumentTypeError("state must be a JSON array containing exactly 14 values") - state: list[float] = [] - for item in raw: - if isinstance(item, bool) or not isinstance(item, int | float) or not math.isfinite(float(item)): - raise argparse.ArgumentTypeError("state values must be finite numbers") - state.append(float(item)) - return state - - -def validate_metadata(metadata: Any) -> dict[str, Any]: - """Validate and summarize the server's advertised action contract.""" - if not isinstance(metadata, dict): - raise ValidationFailure("metadata frame must be a MessagePack object") - expected = { - "protocol_version": PROTOCOL_VERSION, - "robot_profile": "robotwin", - "action_type": ACTION_TYPE, - "action_dim": len(ACTION_ORDER), - "action_dtype": ACTION_DTYPE, - "action_order": list(ACTION_ORDER), - } - for field, expected_value in expected.items(): - if metadata.get(field) != expected_value: - raise ValidationFailure( - f"metadata {field} mismatch: expected {expected_value!r}, observed {metadata.get(field)!r}" - ) - horizon = metadata.get("action_horizon") - if isinstance(horizon, bool) or not isinstance(horizon, int) or horizon < 1: - raise ValidationFailure("metadata action_horizon must be a positive integer") - policy_verified = metadata.get("policy_verified") - verification_status = metadata.get("verification_status") - if not isinstance(policy_verified, bool): - raise ValidationFailure("metadata policy_verified must be boolean") - if not isinstance(verification_status, str) or not verification_status: - raise ValidationFailure("metadata verification_status must be a non-empty string") - max_request_bytes = metadata.get("max_request_bytes") - if isinstance(max_request_bytes, bool) or not isinstance(max_request_bytes, int) or max_request_bytes < 1: - raise ValidationFailure("metadata max_request_bytes must be a positive integer") - scheduling = metadata.get("scheduling") - if not isinstance(scheduling, dict): - raise ValidationFailure("metadata scheduling must be an object") - expected_scheduling = { - "mode": "latest_wins", - "max_pending_per_session": 1, - "sequence_field": "sequence_id", - "ttl_field": "request_ttl_ms", - "inflight_cancellation": False, - } - for field, expected_value in expected_scheduling.items(): - if scheduling.get(field) != expected_value: - raise ValidationFailure( - f"metadata scheduling.{field} mismatch: expected {expected_value!r}, observed {scheduling.get(field)!r}" - ) - max_pending_sessions = scheduling.get("max_pending_sessions") - if isinstance(max_pending_sessions, bool) or not isinstance(max_pending_sessions, int) or max_pending_sessions < 1: - raise ValidationFailure("metadata scheduling.max_pending_sessions must be a positive integer") - return { - "protocol_version": PROTOCOL_VERSION, - "action_shape": [horizon, len(ACTION_ORDER)], - "action_type": ACTION_TYPE, - "action_dtype": ACTION_DTYPE, - "policy_verified": policy_verified, - "verification_status": verification_status, - "max_request_bytes": max_request_bytes, - "scheduling": dict(scheduling), - } - - -def _validate_trace_fields(response: Mapping[str, Any], *, request_id: str, episode_id: str) -> None: - if response.get("request_id") != request_id: - raise ValidationFailure("response did not echo the request_id") - if response.get("episode_id") != episode_id: - raise ValidationFailure("response did not echo the episode_id") - - -def _validate_timings(response: Mapping[str, Any], required_fields: tuple[str, ...]) -> dict[str, float]: - timings = response.get("server_timing") - if not isinstance(timings, dict): - raise ValidationFailure("response server_timing must be an object") - validated: dict[str, float] = {} - for field in required_fields: - value = timings.get(field) - if isinstance(value, bool) or not isinstance(value, int | float) or not math.isfinite(float(value)): - raise ValidationFailure(f"response server_timing.{field} must be a finite number") - if float(value) < 0: - raise ValidationFailure(f"response server_timing.{field} must be non-negative") - validated[field] = float(value) - if "prev_total_ms" in timings: - value = timings["prev_total_ms"] - if isinstance(value, bool) or not isinstance(value, int | float) or not math.isfinite(float(value)): - raise ValidationFailure("response server_timing.prev_total_ms must be a finite number") - if float(value) < 0: - raise ValidationFailure("response server_timing.prev_total_ms must be non-negative") - validated["prev_total_ms"] = float(value) - return validated - - -def validate_reset_response(response: Any, *, request_id: str, episode_id: str) -> dict[str, Any]: - """Validate one reset acknowledgement.""" - if not isinstance(response, dict): - raise ValidationFailure("reset response must be a MessagePack object") - if response.get("action", object()) is not None: - raise ValidationFailure("reset response action must be null") - _validate_trace_fields(response, request_id=request_id, episode_id=episode_id) - return {"server_timing_ms": _validate_timings(response, ("decode_ms", "infer_ms"))} - - -def validate_action_response( - response: Any, - *, - expected_horizon: int, - request_id: str, - episode_id: str, - sequence_id: int, - seed: int, -) -> dict[str, Any]: - """Validate and summarize one raw RoboTwin action response.""" - if not isinstance(response, dict): - raise ValidationFailure("action response must be a MessagePack object") - _validate_trace_fields(response, request_id=request_id, episode_id=episode_id) - if response.get("sequence_id") != sequence_id: - raise ValidationFailure("response did not echo the sequence_id") - if response.get("scheduler_status") != "completed": - raise ValidationFailure("action response scheduler_status must be 'completed'") - if response.get("seed") != seed: - raise ValidationFailure("response did not echo the inference seed") - actions = response.get("action") - if not isinstance(actions, np.ndarray): - raise ValidationFailure("response action must be a NumPy array") - if actions.shape != (expected_horizon, len(ACTION_ORDER)): - raise ValidationFailure( - f"response action shape must be {(expected_horizon, len(ACTION_ORDER))}, got {actions.shape}" - ) - if actions.dtype != np.dtype(np.float32): - raise ValidationFailure(f"response action dtype must be float32, got {actions.dtype}") - if not np.isfinite(actions).all(): - raise ValidationFailure("response action contains non-finite values") - policy_verified = response.get("policy_verified") - verification_status = response.get("verification_status") - if not isinstance(policy_verified, bool): - raise ValidationFailure("response policy_verified must be boolean") - if not isinstance(verification_status, str) or not verification_status: - raise ValidationFailure("response verification_status must be a non-empty string") - - contiguous = np.ascontiguousarray(actions, dtype=" dict[str, Any]: - """Validate a structured scheduler response that must not be executed.""" - if not isinstance(response, dict): - raise ValidationFailure("discarded response must be a MessagePack object") - _validate_trace_fields(response, request_id=request_id, episode_id=episode_id) - if response.get("sequence_id") != sequence_id: - raise ValidationFailure("discarded response did not echo the sequence_id") - if response.get("action", object()) is not None: - raise ValidationFailure("discarded response action must be null") - status = response.get("scheduler_status") - if status not in {"superseded", "expired", "stale_sequence", "overloaded"}: - raise ValidationFailure(f"unexpected scheduler_status: {status!r}") - error = response.get("error") - if not isinstance(error, dict) or error.get("code") != status or not isinstance(error.get("message"), str): - raise ValidationFailure("discarded response must include a matching structured error") - return { - "scheduler_status": status, - "server_timing_ms": _validate_timings(response, ("decode_ms",)), - } - - -def require_exact_replay(records: list[dict[str, Any]]) -> None: - """Require every recorded action digest to match the first response.""" - if len(records) < 2: - raise ValidationFailure("exact replay validation requires at least two requests") - digests = {record["action"]["sha256_float32_le"] for record in records} - if len(digests) != 1: - raise ValidationFailure(f"fixed-seed action replay diverged across {len(digests)} digests") - - -def _receive_binary(connection: Any, *, timeout_seconds: float) -> bytes: - payload = connection.recv(timeout=timeout_seconds) - if not isinstance(payload, bytes): - raise ValidationFailure(f"server returned a text error frame: {payload}") - return payload - - -def load_validation_image(image_path: Path, *, max_image_edge: int) -> tuple[np.ndarray, list[int]]: - """Load an RGB image and bound its encoded request size while preserving aspect ratio.""" - with Image.open(image_path) as opened: - image = opened.convert("RGB") - source_shape = [image.height, image.width, 3] - if max(image.size) > max_image_edge: - image.thumbnail((max_image_edge, max_image_edge), Image.Resampling.LANCZOS) - return np.asarray(image, dtype=np.uint8), source_shape - - -def run_validation( - *, - host: str, - port: int, - image_path: Path, - task: str, - state: list[float], - seed: int, - request_count: int, - timeout_seconds: float, - exact_replay: bool, - request_ttl_ms: float = 5_000.0, - overlap_requests: bool = False, - max_image_edge: int = DEFAULT_MAX_IMAGE_EDGE, -) -> dict[str, Any]: - """Connect to a resident policy and exercise reset plus fixed-seed inference.""" - from websockets.sync.client import connect - - from examples.lingbot_vla_v2.lingbot_vla_v2_robotwin_server import pack_message, unpack_message - - if request_count < 1: - raise ValueError("request_count must be positive") - if exact_replay and request_count < 2: - raise ValueError("exact replay validation requires request_count >= 2") - if overlap_requests and request_count < 2: - raise ValueError("overlap validation requires request_count >= 2") - if overlap_requests and exact_replay: - raise ValueError("overlap validation cannot require exact replay") - if not math.isfinite(request_ttl_ms) or request_ttl_ms <= 0: - raise ValueError("request_ttl_ms must be a positive finite number") - image, source_image_shape = load_validation_image(image_path, max_image_edge=max_image_edge) - - episode_id = f"validator-{time.time_ns()}" - uri = f"ws://{host}:{port}/" - records: list[dict[str, Any]] = [] - started_at = time.perf_counter() - with connect( - uri, - open_timeout=timeout_seconds, - close_timeout=timeout_seconds, - max_size=None, - compression=None, - ) as connection: - metadata = unpack_message(_receive_binary(connection, timeout_seconds=timeout_seconds)) - metadata_summary = validate_metadata(metadata) - expected_horizon = metadata["action_horizon"] - - reset_request_id = "validator-reset" - connection.send( - pack_message( - { - "reset": True, - "robo_name": "robotwin", - "request_id": reset_request_id, - "episode_id": episode_id, - } - ) - ) - reset = unpack_message(_receive_binary(connection, timeout_seconds=timeout_seconds)) - reset_summary = validate_reset_response(reset, request_id=reset_request_id, episode_id=episode_id) - - base_request = { - **{key: image for key in CAMERA_KEYS}, - "observation.state": np.asarray(state, dtype=np.float32), - "task": task, - "episode_id": episode_id, - "seed": seed, - "request_ttl_ms": request_ttl_ms, - } - request_payload_bytes = len( - pack_message( - { - **base_request, - "request_id": "validator-size-check", - "sequence_id": request_count - 1, - } - ) - ) - if request_payload_bytes > metadata["max_request_bytes"]: - raise ValidationFailure( - f"encoded request uses {request_payload_bytes} bytes, exceeding the server limit " - f"of {metadata['max_request_bytes']} bytes" - ) - request_started_at: dict[str, float] = {} - for index in range(request_count): - request_id = f"validator-{index}" - request_started_at[request_id] = time.perf_counter() - connection.send( - pack_message( - { - **base_request, - "request_id": request_id, - "sequence_id": index, - } - ) - ) - if not overlap_requests: - response = unpack_message(_receive_binary(connection, timeout_seconds=timeout_seconds)) - action_summary = validate_action_response( - response, - expected_horizon=expected_horizon, - request_id=request_id, - episode_id=episode_id, - sequence_id=index, - seed=seed, - ) - records.append( - { - "request_id": request_id, - "sequence_id": index, - "round_trip_ms": (time.perf_counter() - request_started_at[request_id]) * 1000.0, - "scheduler_status": "completed", - "action": action_summary, - } - ) - - if overlap_requests: - for _ in range(request_count): - response = unpack_message(_receive_binary(connection, timeout_seconds=timeout_seconds)) - request_id = response.get("request_id") - if not isinstance(request_id, str) or not request_id.startswith("validator-"): - raise ValidationFailure("overlap response has an unknown request_id") - try: - sequence_id = int(request_id.removeprefix("validator-")) - except ValueError as error: - raise ValidationFailure("overlap response request_id has an invalid sequence") from error - if request_id not in request_started_at or not 0 <= sequence_id < request_count: - raise ValidationFailure("overlap response request_id was not sent by this validator") - if any(record["request_id"] == request_id for record in records): - raise ValidationFailure("overlap endpoint returned a duplicate response") - if response.get("scheduler_status") == "completed": - action_summary = validate_action_response( - response, - expected_horizon=expected_horizon, - request_id=request_id, - episode_id=episode_id, - sequence_id=sequence_id, - seed=seed, - ) - record = {"scheduler_status": "completed", "action": action_summary} - else: - discard_summary = validate_discarded_response( - response, - request_id=request_id, - episode_id=episode_id, - sequence_id=sequence_id, - ) - record = {"scheduler_status": discard_summary["scheduler_status"], "discard": discard_summary} - records.append( - { - "request_id": request_id, - "sequence_id": sequence_id, - "round_trip_ms": (time.perf_counter() - request_started_at[request_id]) * 1000.0, - **record, - } - ) - statuses = {record["request_id"]: record["scheduler_status"] for record in records} - if statuses.get(f"validator-{request_count - 1}") != "completed": - raise ValidationFailure("the newest overlap request did not produce an action") - if "superseded" not in statuses.values(): - raise ValidationFailure("overlap validation did not observe a superseded action") - - if exact_replay: - require_exact_replay(records) - return { - "passed": True, - "endpoint": uri, - "elapsed_seconds": time.perf_counter() - started_at, - "configuration": { - "image": str(image_path), - "task": task, - "seed": seed, - "request_count": request_count, - "exact_replay": exact_replay, - "overlap_requests": overlap_requests, - "request_ttl_ms": request_ttl_ms, - "max_image_edge": max_image_edge, - "source_image_shape": source_image_shape, - "transmitted_image_shape": list(image.shape), - "request_payload_bytes": request_payload_bytes, - }, - "metadata": metadata, - "metadata_summary": metadata_summary, - "reset": reset_summary, - "records": records, - } - - -def _positive_int(value: str) -> int: - parsed = int(value) - if parsed < 1: - raise argparse.ArgumentTypeError("value must be positive") - return parsed - - -def _positive_float(value: str) -> float: - parsed = float(value) - if not math.isfinite(parsed) or parsed <= 0: - raise argparse.ArgumentTypeError("value must be a positive finite number") - return parsed - - -def build_parser() -> argparse.ArgumentParser: - """Build the command-line parser.""" - parser = argparse.ArgumentParser(description=__doc__) - parser.add_argument("--host", default="127.0.0.1") - parser.add_argument("--port", type=int, default=9330) - parser.add_argument( - "--image", - type=Path, - default=Path("examples/data/lingbot_world_fast/image.jpg"), - ) - parser.add_argument("--task", default="pick up the object") - parser.add_argument("--state-json", type=parse_state_json, default=[0.0] * 14) - parser.add_argument("--seed", type=int, default=7) - parser.add_argument("--requests", type=_positive_int, default=2) - parser.add_argument("--timeout-seconds", type=_positive_float, default=120.0) - parser.add_argument( - "--request-ttl-ms", - type=_positive_float, - default=5_000.0, - help="Server-local lifetime for each queued action request", - ) - parser.add_argument( - "--overlap-requests", - action="store_true", - help="Send all requests before receiving to verify latest-wins scheduling", - ) - parser.add_argument("--max-image-edge", type=_positive_int, default=DEFAULT_MAX_IMAGE_EDGE) - parser.add_argument("--require-exact-replay", action="store_true") - parser.add_argument("--output", type=Path) - return parser - - -def main() -> None: - """Run the no-simulation RoboTwin WebSocket validation.""" - args = build_parser().parse_args() - report = run_validation( - host=args.host, - port=args.port, - image_path=args.image, - task=args.task, - state=args.state_json, - seed=args.seed, - request_count=args.requests, - timeout_seconds=args.timeout_seconds, - exact_replay=args.require_exact_replay, - request_ttl_ms=args.request_ttl_ms, - overlap_requests=args.overlap_requests, - max_image_edge=args.max_image_edge, - ) - rendered = json.dumps(report, indent=2, sort_keys=True) - if args.output is not None: - args.output.parent.mkdir(parents=True, exist_ok=True) - args.output.write_text(f"{rendered}\n", encoding="utf-8") - print(rendered) - - -if __name__ == "__main__": - main() From a5b7cb8038ddb179cfbe24fb481b7520016d1e27 Mon Sep 17 00:00:00 2001 From: HappyDog0713 Date: Sun, 20 Sep 2026 04:18:39 +0000 Subject: [PATCH 2/4] refactor(vla): remove unused action scheduler Remove the unreferenced ActionChunkScheduler now that session admission and latest-wins handling are owned by the generic VLA WebSocket runtime. Keep the active VLA contracts, replica sessions, simulator execution path, and all LingBot VLA v2 inference modes unchanged. Update VLA documentation to describe the inference-only verification scope, MuJoCo boundary, and retained parity and performance validation. Verification: 130 focused VLA, LingBot, service, replica, and validation tests passed; 237 non-LiveKit service tests passed; runtime import and compile checks passed; git diff --check passed. Full service collection remains blocked only by the optional av dependency, and model-core tests require the pinned Transformers 5.14.1 environment. --- docs/en/vla.md | 35 +++- examples/lingbot_vla_v2/README.md | 5 + telefuser/vla/runtime/__init__.py | 2 - telefuser/vla/runtime/scheduler.py | 270 ----------------------------- 4 files changed, 38 insertions(+), 274 deletions(-) delete mode 100644 telefuser/vla/runtime/scheduler.py diff --git a/docs/en/vla.md b/docs/en/vla.md index f718476..00426b9 100644 --- a/docs/en/vla.md +++ b/docs/en/vla.md @@ -34,7 +34,8 @@ The public modules are: - `telefuser.vla.session`: transport-neutral OPEN, PREDICT, RESET, and CLOSE lifecycle. - `telefuser.vla.serialization`: versioned JSON/Base64 wire formats for action and observation spaces, observations, and chunks. -- `telefuser.vla.runtime`: scheduling, deterministic chunk state, action trimming, age checks, and safety policies. +- `telefuser.vla.runtime`: deterministic chunk state, action trimming, age checks, safety policies, and simulator + execution reports. Session admission and latest-wins request handling live in `telefuser.service.vla_session`. - `telefuser.integrations.sim`: simulator protocol and the dependency-free RoboTwin callback adapter. - `telefuser.service.vla_session`: the additive generic VLA WebSocket application factory. - `telefuser.service.vla_replica`: worker-local OPEN, PREDICT, RESET, and CLOSE dispatch for pipeline replicas. @@ -53,6 +54,11 @@ The generic VLA WebSocket is the online integration path for both RoboTwin and M an explicit model and embodiment pair, while the shared HTTP structured-task route remains unchanged and continues to return the existing canonical JSON result. +LingBot-VLA v2 exposes three access modes over the same pipeline: `lingbot_vla_v2_inference.py` is the direct offline +baseline, `lingbot_vla_v2_native_service.py` is the existing structured HTTP compatibility path, and +`lingbot_vla_v2_vla_server.py` is the continuous-control WebSocket path. Only the WebSocket path returns the decoded +semantic `RobotActionChunk`; the direct and native paths retain the canonical normalized model result. + ## Session Lifecycle Create and register loaded components, then open a typed session: @@ -131,9 +137,34 @@ claim that a returned action was executed. `SimulatorChunkRuntime` runs on the s EXECUTING, EXECUTED, HOLD, and STOP transitions while calling a `SimulatorAdapter` one action at a time. The bundled RoboTwin statistics do not declare an authoritative simulator control frequency, so both exported -LingBot/RoboTwin specs use `control_hz=None`. The remote simulator must resolve its actual control period rather than +LingBot/RobotWin specs use `control_hz=None`. The remote simulator must resolve its actual control period rather than guessing one in the inference process. +## LingBot-VLA v2 Verification Scope + +The TeleFuser integration consumes the official checkpoint as-is. It does not add post-training, fine-tuning, or +model-specific camera-pose inputs. The supported verification path is: + +```text +direct pipeline baseline + -> native structured HTTP compatibility + -> generic VLA WebSocket session + -> simulator-side RobotWin execution +``` + +The direct and native paths return the canonical normalized `50 x 55` action chunk. The generic session path applies +`RobotWinProfile` and returns an absolute-position `50 x 14` `RobotActionChunk`. MuJoCo exercises the simulator-side +execution boundary; it does not change model computation or action semantics. + +Keep inference speed comparisons and release artifacts separate from task-success claims. The maintained validation +tools cover upstream tensor parity, runtime latency, quantization comparisons, structured-service behavior, replica +faults, and shutdown/restart. Generated captures and reports belong under the ignored `work_dirs/` directory. A +real-robot task-success evaluation is outside this integration and must not be represented by `policy_verified`. + +The generic WebSocket remains an explicit standalone service so existing `telefuser serve` routes and non-VLA +pipelines stay unchanged. Deployment authentication, TLS, actuator feedback, and emergency-stop behavior remain +outside the inference process and must be supplied by the deployment or simulator boundary. + ## Simulator Boundary `RoboTwinSimulatorAdapter` has no RoboTwin, SAPIEN, Vulkan, or ROS dependency. The RTX-side process provides observe, diff --git a/examples/lingbot_vla_v2/README.md b/examples/lingbot_vla_v2/README.md index 3ef0a75..11ed30b 100644 --- a/examples/lingbot_vla_v2/README.md +++ b/examples/lingbot_vla_v2/README.md @@ -309,6 +309,11 @@ server. This is a chain/physics smoke test, not a RoboTwin task-success or check The repository includes strict upstream parity, runtime, quantization, structured-service, fault, and AIPerf validators under `tools/validation/` and `benchmarks/telefuser_aiperf/`. +This integration is inference-only: it consumes the official LingBot-VLA v2 checkpoint without post-training or +fine-tuning. The MuJoCo example validates the simulator boundary and action execution path; it is not a substitute +for physical robot task-success evaluation. Keep direct/HTTP/WebSocket speed comparisons and their generated reports +under the ignored `work_dirs/` directory. + Compare previously captured upstream and TeleFuser artifacts: ```bash diff --git a/telefuser/vla/runtime/__init__.py b/telefuser/vla/runtime/__init__.py index 7ce4c57..f7785a9 100644 --- a/telefuser/vla/runtime/__init__.py +++ b/telefuser/vla/runtime/__init__.py @@ -10,13 +10,11 @@ ) from .executor import ChunkExecutor, RobotAction from .safety import ActionSafetyPolicy, BoundedActionSafety, FiniteActionSafety -from .scheduler import ActionChunkScheduler from .simulator import ChunkExecutionReport, SimulatorChunkRuntime __all__ = [ "ActionChunkStateMachine", "ActionSafetyPolicy", - "ActionChunkScheduler", "BoundedActionSafety", "ChunkExecutor", "ChunkExecutionReport", diff --git a/telefuser/vla/runtime/scheduler.py b/telefuser/vla/runtime/scheduler.py deleted file mode 100644 index 24669f6..0000000 --- a/telefuser/vla/runtime/scheduler.py +++ /dev/null @@ -1,270 +0,0 @@ -"""Bounded latest-wins scheduling for action chunk inference.""" - -from __future__ import annotations - -import asyncio -import math -import operator -import time -from collections import OrderedDict -from dataclasses import dataclass -from typing import Any, Callable, Mapping - - -@dataclass(slots=True) -class _ScheduledAction: - request: Mapping[str, Any] - session_key: str - generation: int - received_at: float - deadline_at: float | None - future: asyncio.Future[dict[str, Any]] - - -class ActionChunkScheduler: - """Serialize inference while retaining only the newest pending chunk per session.""" - - def __init__( - self, - infer: Callable[[Mapping[str, Any]], dict[str, Any]], - *, - max_pending_sessions: int = 32, - ) -> None: - if max_pending_sessions < 1: - raise ValueError("max_pending_sessions must be positive") - self._infer = infer - self._max_pending_sessions = max_pending_sessions - self._pending: OrderedDict[str, _ScheduledAction] = OrderedDict() - self._latest_generation: dict[str, int] = {} - self._latest_sequence: dict[str, int] = {} - self._wake = asyncio.Event() - self._worker: asyncio.Task[None] | None = None - self._closed = False - - @property - def metadata(self) -> dict[str, Any]: - """Describe the additive scheduling controls accepted by the endpoint.""" - return { - "scheduling": { - "mode": "latest_wins", - "max_pending_per_session": 1, - "max_pending_sessions": self._max_pending_sessions, - "sequence_field": "sequence_id", - "ttl_field": "request_ttl_ms", - "inflight_cancellation": False, - } - } - - async def start(self) -> None: - """Start the single inference worker on the current event loop.""" - if self._worker is not None: - return - if self._closed: - raise RuntimeError("action scheduler is closed") - self._worker = asyncio.create_task(self._run(), name="vla-action-chunk-scheduler") - - def submit(self, request: Mapping[str, Any], *, session_key: str) -> asyncio.Future[dict[str, Any]]: - """Admit one request and return a future without waiting for inference.""" - if self._worker is None or self._closed: - raise RuntimeError("action scheduler is not running") - if not session_key: - raise ValueError("session_key must be non-empty") - - loop = asyncio.get_running_loop() - future: asyncio.Future[dict[str, Any]] = loop.create_future() - received_at = time.monotonic() - ttl_ms = self._optional_ttl_ms(request) - sequence_id = self._optional_sequence_id(request) - previous = self._pending.get(session_key) - if previous is None and len(self._pending) >= self._max_pending_sessions: - future.set_result( - self._discarded_response( - request, - status="overloaded", - message="the action scheduler has no free pending-session slot", - ) - ) - return future - if sequence_id is not None: - latest_sequence = self._latest_sequence.get(session_key) - if latest_sequence is not None and sequence_id <= latest_sequence: - future.set_result( - self._discarded_response( - request, - status="stale_sequence", - message=f"sequence_id={sequence_id} is not newer than {latest_sequence}", - ) - ) - return future - self._latest_sequence[session_key] = sequence_id - - generation = self._latest_generation.get(session_key, 0) + 1 - self._latest_generation[session_key] = generation - deadline_at = None if ttl_ms is None else received_at + ttl_ms / 1000.0 - job = _ScheduledAction( - request=dict(request), - session_key=session_key, - generation=generation, - received_at=received_at, - deadline_at=deadline_at, - future=future, - ) - if previous is not None: - self._resolve( - previous, - self._discarded_response( - previous.request, - status="superseded", - message="a newer observation replaced this pending action request", - ), - ) - self._pending[session_key] = job - self._wake.set() - return future - - def release_session(self, session_key: str) -> None: - """Discard pending work and sequence state for a disconnected session.""" - pending = self._pending.pop(session_key, None) - if pending is not None: - self._resolve( - pending, - self._discarded_response( - pending.request, - status="session_closed", - message="the client session closed before inference", - ), - ) - self._latest_generation.pop(session_key, None) - self._latest_sequence.pop(session_key, None) - - async def close(self) -> None: - """Drain an in-flight call and reject work that has not started.""" - if self._closed: - return - self._closed = True - for job in self._pending.values(): - self._resolve( - job, - self._discarded_response( - job.request, - status="server_stopping", - message="the action scheduler is stopping", - ), - ) - self._pending.clear() - self._wake.set() - if self._worker is not None: - await self._worker - self._worker = None - - async def _run(self) -> None: - while True: - await self._wake.wait() - if self._closed and not self._pending: - return - if not self._pending: - self._wake.clear() - continue - - _, job = self._pending.popitem(last=False) - if not self._pending: - self._wake.clear() - if job.future.done(): - continue - discard = self._discard_reason(job) - if discard is not None: - self._resolve(job, discard) - continue - - inference_started_at = time.monotonic() - try: - response = await asyncio.to_thread(self._infer, job.request) - except Exception as error: - if not job.future.done(): - job.future.set_exception(error) - continue - inference_ms = (time.monotonic() - inference_started_at) * 1000.0 - - discard = self._discard_reason(job) - if discard is not None: - self._resolve(job, discard) - continue - completed_at = time.monotonic() - response = dict(response) - server_timing = dict(response.get("server_timing", {})) - server_timing.update( - queue_wait_ms=(inference_started_at - job.received_at) * 1000.0, - infer_ms=inference_ms, - scheduler_total_ms=(completed_at - job.received_at) * 1000.0, - ) - response.update(scheduler_status="completed", server_timing=server_timing) - for field in ("request_id", "episode_id", "sequence_id"): - if field in job.request: - response.setdefault(field, job.request[field]) - if job.deadline_at is not None: - response["request_ttl_ms"] = (job.deadline_at - job.received_at) * 1000.0 - self._resolve(job, response) - - def _discard_reason(self, job: _ScheduledAction) -> dict[str, Any] | None: - if job.deadline_at is not None and time.monotonic() >= job.deadline_at: - return self._discarded_response( - job.request, - status="expired", - message="the action request exceeded request_ttl_ms", - ) - if self._latest_generation.get(job.session_key) != job.generation: - return self._discarded_response( - job.request, - status="superseded", - message="a newer observation superseded this action result", - ) - return None - - @staticmethod - def _resolve(job: _ScheduledAction, response: dict[str, Any]) -> None: - if not job.future.done(): - job.future.set_result(response) - - @staticmethod - def _optional_sequence_id(request: Mapping[str, Any]) -> int | None: - value = request.get("sequence_id") - if value is None: - return None - if isinstance(value, bool): - raise ValueError("sequence_id must be a non-negative integer") - try: - sequence_id = operator.index(value) - except TypeError as error: - raise ValueError("sequence_id must be a non-negative integer") from error - if sequence_id < 0: - raise ValueError("sequence_id must be a non-negative integer") - return sequence_id - - @staticmethod - def _optional_ttl_ms(request: Mapping[str, Any]) -> float | None: - value = request.get("request_ttl_ms") - if value is None: - return None - if isinstance(value, bool) or not isinstance(value, int | float): - raise ValueError("request_ttl_ms must be a positive finite number") - ttl_ms = float(value) - if not math.isfinite(ttl_ms) or ttl_ms <= 0: - raise ValueError("request_ttl_ms must be a positive finite number") - return ttl_ms - - @staticmethod - def _discarded_response( - request: Mapping[str, Any], - *, - status: str, - message: str, - ) -> dict[str, Any]: - response: dict[str, Any] = { - "action": None, - "scheduler_status": status, - "error": {"code": status, "message": message}, - } - for field in ("request_id", "episode_id", "sequence_id"): - if field in request: - response[field] = request[field] - return response From 9ca1c7c58f1ced4cb49adee61501b693e35cafe9 Mon Sep 17 00:00:00 2001 From: HappyDog0713 Date: Sun, 20 Sep 2026 07:32:11 +0000 Subject: [PATCH 3/4] docs: add VLA serving milestones to news Document the general VLA session-serving abstraction and the LingBot-VLA v2 integration milestone in the English and Chinese homepage News sections.\n\nHighlight semantic observation/action contracts, extensible embodiment adapters, replica-affine sessions, WebSocket continuous-control serving, MuJoCo validation, structured HTTP, BF16 eager, CUDA Graph, and quantized inference support.\n\nVerification:\n- git diff --check --- README.md | 6 ++++++ README_zh.md | 4 ++++ 2 files changed, 10 insertions(+) diff --git a/README.md b/README.md index b666396..54c7099 100644 --- a/README.md +++ b/README.md @@ -25,6 +25,12 @@ runtime path, supported workloads, and reproducible real-time gate. ## News 📰 +- ✨ **2026-09-14**: Established a general VLA session-serving abstraction with validated semantic observation and + action contracts, extensible embodiment adapters, replica-affine sessions, continuous-control WebSocket serving, + and simulator-side MuJoCo validation. +- ✨ **2026-08-27**: Added **LingBot-VLA v2** 6B base-checkpoint support with RobotWin multi-camera, task-text, + robot-state, and normalized action I/O, together with native structured HTTP serving, BF16 eager, CUDA Graph, + and quantized inference support. - ✨ **2026-08-19**: Added [**LTX-2.5 Distilled**](examples/ltx25_distilled/README.md) T2V and I2V joint audio-video generation with a ModuleManager-backed six-stage pipeline, selectable dense attention backends, and Ulysses sequence parallelism on **1, 2, or 4 x H100** GPUs. diff --git a/README_zh.md b/README_zh.md index 4ef492c..ce9ea92 100644 --- a/README_zh.md +++ b/README_zh.md @@ -24,6 +24,10 @@ TeleFuser 是一个开源的多模态生成与世界模型流式推理和服务 ## News 📰 +- ✨ **2026-09-14**:建立通用 VLA 会话服务抽象,定义可校验的语义观测与动作契约, + 支持可扩展具身适配器、模型副本绑定会话、连续控制 WebSocket 服务及模拟器侧 MuJoCo 验证。 +- ✨ **2026-08-27**:接入 **LingBot-VLA v2** 6B base checkpoint,支持 RobotWin 多相机、任务文本、 + 机器人状态和归一化动作输入输出,并提供原生结构化 HTTP、BF16 eager、CUDA Graph 及量化推理支持。 - ✨ **2026-08-19**:新增 [**LTX-2.5 Distilled**](examples/ltx25_distilled/README.md) T2V 和 I2V 联合 音视频生成,采用基于 ModuleManager 的六阶段 Pipeline,支持选择密集注意力后端,并可在 **1、2 或 4 张 H100** 上使用 Ulysses 序列并行。 From daec2ff0e25481500809f2fcf40c20a84eb18d89 Mon Sep 17 00:00:00 2001 From: HappyDog0713 Date: Sun, 20 Sep 2026 08:07:18 +0000 Subject: [PATCH 4/4] docs(vla): publish optimization performance comparisons Link the LingBot-VLA v2 News entry to its detailed guide and document CUDA Graph and quantization profile performance comparisons.\n\nAdd H100 BF16 eager versus denoising Graph A/B measurements, quantization latency and VRAM tables, benchmark scope, and release-validation caveats. Keep quantized profiles clearly described as capacity and memory trade-offs where exact replay validation is incomplete.\n\nVerification:\n- git diff --check --- README.md | 6 ++--- README_zh.md | 4 +-- examples/lingbot_vla_v2/README.md | 41 +++++++++++++++++++++++++++++++ 3 files changed, 46 insertions(+), 5 deletions(-) diff --git a/README.md b/README.md index 54c7099..3c35b0a 100644 --- a/README.md +++ b/README.md @@ -28,9 +28,9 @@ runtime path, supported workloads, and reproducible real-time gate. - ✨ **2026-09-14**: Established a general VLA session-serving abstraction with validated semantic observation and action contracts, extensible embodiment adapters, replica-affine sessions, continuous-control WebSocket serving, and simulator-side MuJoCo validation. -- ✨ **2026-08-27**: Added **LingBot-VLA v2** 6B base-checkpoint support with RobotWin multi-camera, task-text, - robot-state, and normalized action I/O, together with native structured HTTP serving, BF16 eager, CUDA Graph, - and quantized inference support. +- ✨ **2026-08-27**: Added [**LingBot-VLA v2**](examples/lingbot_vla_v2/README.md) 6B base-checkpoint support + with RobotWin multi-camera, task-text, robot-state, and normalized action I/O, together with native structured HTTP + serving, BF16 eager, CUDA Graph, and quantized inference support. - ✨ **2026-08-19**: Added [**LTX-2.5 Distilled**](examples/ltx25_distilled/README.md) T2V and I2V joint audio-video generation with a ModuleManager-backed six-stage pipeline, selectable dense attention backends, and Ulysses sequence parallelism on **1, 2, or 4 x H100** GPUs. diff --git a/README_zh.md b/README_zh.md index ce9ea92..05ad6b0 100644 --- a/README_zh.md +++ b/README_zh.md @@ -26,8 +26,8 @@ TeleFuser 是一个开源的多模态生成与世界模型流式推理和服务 - ✨ **2026-09-14**:建立通用 VLA 会话服务抽象,定义可校验的语义观测与动作契约, 支持可扩展具身适配器、模型副本绑定会话、连续控制 WebSocket 服务及模拟器侧 MuJoCo 验证。 -- ✨ **2026-08-27**:接入 **LingBot-VLA v2** 6B base checkpoint,支持 RobotWin 多相机、任务文本、 - 机器人状态和归一化动作输入输出,并提供原生结构化 HTTP、BF16 eager、CUDA Graph 及量化推理支持。 +- ✨ **2026-08-27**:接入 [**LingBot-VLA v2**](examples/lingbot_vla_v2/README.md) 6B base checkpoint,支持 RobotWin 多相机、 + 任务文本、机器人状态和归一化动作输入输出,并提供原生结构化 HTTP、BF16 eager、CUDA Graph 及量化推理支持。 - ✨ **2026-08-19**:新增 [**LTX-2.5 Distilled**](examples/ltx25_distilled/README.md) T2V 和 I2V 联合 音视频生成,采用基于 ModuleManager 的六阶段 Pipeline,支持选择密集注意力后端,并可在 **1、2 或 4 张 H100** 上使用 Ulysses 序列并行。 diff --git a/examples/lingbot_vla_v2/README.md b/examples/lingbot_vla_v2/README.md index 11ed30b..e79ebcf 100644 --- a/examples/lingbot_vla_v2/README.md +++ b/examples/lingbot_vla_v2/README.md @@ -189,6 +189,47 @@ Compare a deterministic quantized capture with the corresponding TeleFuser BF16 --output work_dirs/vla_quantization/bf16_vs_torchao.json ``` +## Performance + +The measurements below are point results from one NVIDIA H100 80 GB system with CUDA 13.0, PyTorch 2.11.0+cu130, +Transformers 5.14.1, batch size 1, and fixed seed 7. Core-model timings reuse device-resident parity inputs and +exclude image decoding and preprocessing. They are environment-specific measurements, not universal performance +guarantees. + +### CUDA Graph A/B + +The graph captures only the fixed ten-step action-denoising loop; vision-language prefix encoding remains eager. The +comparison used five warmup requests and twenty serial requests with the same SDPA configuration in both paths. + +| Scope | BF16 eager | BF16 denoising Graph | Change | Speedup | +| --- | ---: | ---: | ---: | ---: | +| Direct core model | 649.790 ms | 164.226 ms | -74.73% | 3.957x | +| Direct runtime request | 649.843 ms | 164.910 ms | -74.62% | 3.941x | +| Service target inference | 1317.571 ms | 883.986 ms | -32.91% | 1.490x | +| Service HTTP end-to-end | 1362.134 ms | 929.761 ms | -31.74% | 1.465x | + +### Quantization profiles + +The table reports core-model latency and peak allocated VRAM from the H100 release profiles. The relative value is +`BF16 eager latency / profile latency`; values below `1.0x` are slower than BF16 eager. `fused-fp8-graph` combines +quantization with CUDA Graph and must be compared with BF16 Graph when isolating the quantization effect. + +| Profile | Mean core latency | Relative to BF16 eager | Peak allocated VRAM | Interpretation | +| --- | ---: | ---: | ---: | --- | +| BF16 eager | 637.7 ms | 1.00x | 12,457 MiB | Reference | +| BF16 CUDA Graph | 131.6 ms | 4.84x | 12,487 MiB | Graph-only reference | +| `fused-fp8-graph` | 162.5 ms | 3.92x | 10,862 MiB | Combined FP8 + Graph path; 0.81x vs BF16 Graph | +| `torchao-fp8` | 1335.0 ms | 0.48x | 8,438 MiB | Lower memory, slower than BF16 eager | +| `bnb-nf4` | 924.7 ms | 0.69x | 6,479 MiB | Lower memory, slower than BF16 eager | + +The CUDA Graph A/B table and the release-profile table use separate benchmark harnesses and run sets; compare +absolute timings only within the same table. + +The quantized profiles are capacity and memory trade-off profiles rather than release-validated speedup claims. The +runnable profiles passed functional, numerical-threshold, AIPerf, fault, and shutdown checks, but did not produce +bit-exact HTTP replays. The `tf-kernel-fp8` profile is excluded from this table because compatible hardware validation +is still pending. Re-run the release suite to regenerate measurements for a different GPU or software environment. + ## Serving The three inference entrypoints share one `LingBotVlaV2Pipeline`; they are access modes, not separate model