Skip to content

Support DROID policy server for Cosmos3 OpenPI - #4282

Merged
hsliuustc0106 merged 12 commits into
vllm-project:mainfrom
yuzhudong:yuzhud/mbala/cosmos3_action_review
Jun 15, 2026
Merged

hsliuustc0106 merged 12 commits into
vllm-project:mainfrom
yuzhudong:yuzhud/mbala/cosmos3_action_review

Conversation

@yuzhudong

@yuzhudong yuzhudong commented Jun 9, 2026 •

Copy link
Copy Markdown
Contributor

Purpose

  • Support the DROID policy server/OpenPI action-serving flow for Cosmos3.
  • Add parity bridge updates for Cosmos3 action outputs.
  • Align OpenPI connection/serving tests with the review feedback.

Test Plan

  • Syntax-check the touched Python files.
  • Existing OpenPI connection/serving tests cover the updated behavior.

Test Result

python -m py_compile vllm_omni/diffusion/models/cosmos3/pipeline_cosmos3.py vllm_omni/diffusion/models/cosmos3/transformer_cosmos3.py vllm_omni/diffusion/diffusion_engine.py vllm_omni/entrypoints/openpi/connection.py tests/entrypoints/openai_api/test_openpi_connection.py tests/entrypoints/openai_api/test_openpi_serving.py

Passed.

@chatgpt-codex-connector

Copy link
Copy Markdown

Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits.
Credits must be used to enable repository wide code reviews.

@yuzhudong
yuzhudong force-pushed the yuzhud/mbala/cosmos3_action_review branch from d08f34c to 0ba7b37 Compare June 9, 2026 05:43

@hsliuustc0106 hsliuustc0106 left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cosmos3 DROID policy server support for OpenPI looks correct. Parity bridge updates and test alignment are appropriate.

@david6666666

Copy link
Copy Markdown
Collaborator

Can you provide specific end-to-end test videos or results for the DROID robot platform?

@lishunyang12 lishunyang12 left a comment •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

a few comments about code quality and please fix DCO

nd = _mapping_get(obj, "nd")
dtype = _mapping_get(obj, "type")
data = _mapping_get(obj, "data")
if nd is not None and dtype is not None and data is not None:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Marker detection is weaker than the vendored codec it replaces: any user dict carrying nd/type/data keys is silently decoded as a numpy array. Reserve a unique marker (or also require kind) and document the wire contract.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: unpack now requires the kind marker, validates it against the dtype, and the OpenPI msgpack-numpy wire contract is documented in the module docstring.

dtype = _mapping_get(obj, "type")
data = _mapping_get(obj, "data")
if nd is not None and dtype is not None and data is not None:
array = np.frombuffer(data, dtype=np.dtype(dtype))

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

np.frombuffer returns a read-only array. Downstream RoboLab code mutates decoded arrays in place (action_np[:, -1] = ...) and will raise. Return .copy().

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: _unpack_numpy() now calls .copy() on the np.frombuffer(...) result before reshape/scalar extraction, so downstream mutation is safe.

b"nd": True,
b"data": obj.tobytes(),
b"type": obj.dtype.str,
b"kind": obj.dtype.kind,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

b"kind" is written on pack but never read on unpack. Drop it or use it — dead wire fields invite drift.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: kind is now read on unpack and validated against np.dtype(...).kind, so it is no longer a dead wire field.

Comment thread vllm_omni/diffusion/diffusion_engine.py Outdated

custom_output = output.custom_output or {}
action_payload = custom_output.get("actions")
action_only_output = action_payload is not None and isinstance(output_data, dict) and not output_data

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

"Empty output_data dict means action-only" is an implicit cross-file contract. If any other path returns {}, post-processing is silently skipped. Prefer an explicit flag in custom_output.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863/c7313679: RoboLab action-only output is now explicit via custom_output['action_only_output'], and both DiffusionEngine and video serving consume that flag instead of relying on an empty output_data dict.

model_audio_sample_rate = outputs.get("audio_sample_rate")
model_fps = outputs.get("fps")
outputs = outputs.get("video", outputs)
if action_payload is None:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Dead code: action_payload was already assigned from custom_output.get("actions") above, so this re-fetch can never change the value. Remove.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in 0763d5aa: action_payload is now populated from an output dict only when one is present, then falls back to custom_output['actions']. This removes the self-referential outputs.get('actions', action_payload) fallback while preserving both output paths.

domain_name = str(extra_param("domain_name", ROBOLAB_DEFAULT_DOMAIN_NAME))
domain_id = resolve_domain_id(domain_name=domain_name, require_explicit=True)

if history_length < (1 if use_state else 0):

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Message says ">= 1 when use_state is true" but the else 0 branch is always satisfied, so the text doesn't match the check. Tighten.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: tightened the validation to if use_state and history_length < 1 so the condition and error message match.

"ai_caption": prompt,
"video": video,
"action": action,
"conditioning_fps": torch.tensor(fps, dtype=torch.long),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

fps is a float (default 15.0) but torch.long truncates fractional fps. Intentional for parity? If so, comment it.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: kept the torch.long conversion for Cosmos Framework parity and added a comment documenting that it is consumed as an integer conditioning bucket.

else:
transformed_raw_action_dim = raw_action_dim

return {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Returning a ~20-key dict[str, Any] accessed as robolab_inputs["..."] across forward() is stringly-typed and typo-prone. A @dataclass makes the contract explicit and checkable.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: introduced the frozen RoboLabPolicyInputs dataclass and replaced the string-key access through forward() with typed attributes.

initial_pose[:3, :3] = _convert_midtrain_rotation(eef_quat[-1], "quat_xyzw", "matrix")
initial_pose[:3, 3] = eef_pos[-1]
abs_pose = _pose_rel_to_abs(
action_np[:, :9],

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Magic 9 (and :9/9: slices) = pos(3)+rot6d(6), knowable only to the author. Name the constant.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f15c9863: added ROBOLAB_MIDTRAIN_POSE_ACTION_DIM and replaced the :9 / 9: slices with that named constant.

sp = req.sampling_params
prompt_data = req.prompts[0]
if isinstance(prompt_data, str):
robolab_inputs = self._build_robolab_policy_inputs(sp)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

~12 scattered if robolab_inputs is not None: forks thread a second code path through an already-huge forward(), sharply raising complexity. Extract the RoboLab flow into its own method/strategy.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f91acdfe: extracted the RoboLab policy path into _forward_robolab_policy(). forward() now dispatches once for RoboLab requests, and the normal generation path no longer threads the RoboLab-specific branches.

@lishunyang12

lishunyang12 commented Jun 10, 2026 •

Copy link
Copy Markdown
Collaborator

Also please consider update recipe under vllm-omni/receipes/cosmos3.

@yuzhudong
yuzhudong force-pushed the yuzhud/mbala/cosmos3_action_review branch from 3e07645 to c731367 Compare June 10, 2026 17:47
@yuzhudong

Copy link
Copy Markdown
Contributor Author

@lishunyang12 addressed the Cosmos3 recipe update in f15c9863: recipes/cosmos3/Cosmos3-Nano.md now includes the DROID OpenPI policy server flow plus the observation and extra-param notes.

@yuzhudong
yuzhudong force-pushed the yuzhud/mbala/cosmos3_action_review branch from b5bc127 to d26fb1b Compare June 10, 2026 22:17
return keep


def _ensure_rgb_uint8_image(value: Any, key: str) -> np.ndarray:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The clutter from all these local functions is getting quite large, even though these are not necessarily crucial to understanding Cosmos3 model and pipeline. Can you refactor the code by moving all the local functions like that to vllm_omni/diffusion/models/cosmos3/utils.py?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f93d71c3: moved the Cosmos3/RoboLab helper constants, dataclass, image/action utilities, lazy loaders, and RoboLab action conversion into vllm_omni/diffusion/models/cosmos3/utils.py. pipeline_cosmos3.py now imports those helpers and keeps the main pipeline flow focused on request construction and denoising.

Comment on lines -258 to +264
if self.post_process_func is not None:
if action_only_output:
outputs = []
elif self.post_process_func is not None:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I see we still use action-specific post-processing function _postprocess_robolab_action, just inside the Cosmos3 pipeline instead of here. I think we can allow for framework-level post-processing function for action output as well and refactor the code so that _postprocess_robolab_action is assigned as this action-specific post-processing function and called here, instead of inside the pipeline

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in f93d71c3: added a registered diffusion action postprocess hook and wired DiffusionEngine.step() to call it for action payloads. Cosmos3 now registers get_cosmos3_action_post_process_func; the RoboLab pipeline returns raw action output plus RoboLabPolicyInputs, and the engine-level action hook runs the RoboLab-specific conversion.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up in c879db5b: the closed-loop RoboLab eval regressed with the conversion moved fully across the pipeline/engine boundary, so this preserves the previous RoboLab outward action path: the pipeline again emits final actions and no internal robolab_policy_inputs object. The generic engine action postprocess hook remains, but it now only runs when a model has not already supplied final actions, so it cannot clobber model-provided action payloads.

@MaciejBalaNV

Copy link
Copy Markdown
Contributor

LGTM now, thanks, there's a conflict to be resolved though after a recent PR for guardrail exception was merged.

@yuzhudong
yuzhudong force-pushed the yuzhud/mbala/cosmos3_action_review branch from e935bbb to 464040c Compare June 12, 2026 17:50
@yuzhudong
yuzhudong requested a review from ywang96 as a code owner June 12, 2026 17:50
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: yuzhud <yuzhudong@gmail.com>
This reverts commit d47d257.

Signed-off-by: yuzhud <yuzhudong@gmail.com>
@yuzhudong
yuzhudong force-pushed the yuzhud/mbala/cosmos3_action_review branch from 464040c to 54e0ab1 Compare June 12, 2026 17:52
@yuzhudong

Copy link
Copy Markdown
Contributor Author

LGTM now, thanks, there's a conflict to be resolved though after a recent PR for guardrail exception was merged.

Rebased to the latest main

@hsliuustc0106 hsliuustc0106 added the ready label to trigger buildkite CI label Jun 13, 2026
)


def _uses_cosmos3_diffusion_model(od_config: OmniDiffusionConfig, model_config: _DiffusionVllmModelConfig) -> bool:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

can we move the model specific checks outside the diffusion worker?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Addressed in 0df9a175: removed the Cosmos3 model check from DiffusionWorker. The worker now asks the diffusion registry for an optional model IR-op-priority hook instead of carrying model-family-specific logic itself.

@hsliuustc0106 hsliuustc0106 left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Posted focused inline comments on the diffusion engine/worker scalability concerns from the review.

# Also need to log, because vLLM internally logs another line in VllmConfig.__post_init__. Avoid confusion.
vllm_config.kernel_config.ir_op_priority = current_omni_platform.get_default_ir_op_priority(vllm_config)
ir_op_priority = current_omni_platform.get_default_ir_op_priority(vllm_config)
if _uses_cosmos3_diffusion_model(self.od_config, vllm_config.model_config):

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This puts Cosmos3-specific kernel policy directly in the generic diffusion worker. That makes DiffusionWorker.init_device() harder to scale as more model families need op-priority tweaks, and it also replaces the platform default wholesale. Can this be expressed as a model/platform config hook, then merged with current_omni_platform.get_default_ir_op_priority(...) instead of special-casing Cosmos3 here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Makes sense; addressed in 0df9a175. DiffusionWorker now starts from current_omni_platform.get_default_ir_op_priority(vllm_config) and passes that default into a registered model hook. The Cosmos3 hook copies the platform priority fields and overrides only rms_norm / fused_add_rms_norm to native, so other platform defaults are preserved instead of replacing the whole config.

if action_payload is None:
action_payload = custom_output.get("actions")
action_post_process_func = getattr(self, "action_post_process_func", None)
if action_payload is None and action_post_process_func is not None:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new action postprocess hook is only used when custom_output["actions"] is absent, but the Cosmos3 policy path already writes processed actions before the engine sees it. That leaves action postprocessing split between pipeline code and engine hooks, and the registered hook is bypassed for the only current user. I’d prefer one owner: either pipelines return raw action and the engine hook always owns action postprocess, or this stays inside the normal model postprocess path without adding a separate generic hook.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed; addressed in 0df9a175. The RoboLab policy path no longer writes final processed actions in the pipeline. It returns raw action plus compact RoboLabActionPostprocessInputs, and the registered Cosmos3 action hook performs the RoboLab conversion in DiffusionEngine.step(). I kept the metadata small so we do not send the full RoboLab video/input object across the worker boundary.


postprocess_start_time = time.perf_counter()
if self.post_process_func is not None:
if action_only_output:

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This adds another modality-specific branch to DiffusionEngine.step, which is already doing image/audio/text response shaping and multibatch slicing. It works for this PR, but it makes the engine the long-term router for every new output type. Consider introducing a typed postprocess/output envelope so modality-specific payload extraction and slicing can live outside the engine core.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I agree with the direction, but I did not fold a typed output-envelope redesign into this PR. That would be a broader cross-model engine/API change because step() currently handles existing image, audio, text, and video shaping conventions. For this PR I kept the change scoped to the concrete ownership problem: action conversion is now owned by the registered model action hook, and the worker-specific Cosmos3 kernel policy was moved behind a registered model hook. A typed output envelope would be a good follow-up cleanup once the generic diffusion output contract is designed across modalities.

Signed-off-by: yuzhud <yuzhudong@gmail.com>
@hsliuustc0106 hsliuustc0106 added ready label to trigger buildkite CI and removed ready label to trigger buildkite CI labels Jun 13, 2026
@hsliuustc0106
hsliuustc0106 merged commit e2e7917 into vllm-project:main Jun 15, 2026
7 of 8 checks passed
AbelSara pushed a commit to AbelSara/vllm-omni that referenced this pull request Jun 16, 2026
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Signed-off-by: z00806815 <zhuhonghan@huawei.com>
loveysuby added a commit to loveysuby/vllm-omni that referenced this pull request Jun 17, 2026
The robot serving path (/v1/realtime/robot/openpi) was broken for real
openpi-client usage. connection.py packed/unpacked numpy arrays with the
standard msgpack-numpy markers ({nd, type, kind}) introduced in vllm-project#4282,
but openpi-client uses its own markers ({__ndarray__, __npgeneric__}).
Neither direction could decode the other, so every inference failed with
"Internal inference error" (np.object_ array -> torch.from_numpy TypeError
in transform/droid.py:_preprocess_view).

- _pack_numpy now emits the openpi-client wire format so clients can
  decode server responses.
- _unpack_numpy accepts both the openpi-client markers and the legacy
  {nd, type, kind} markers for backward compatibility.
- Add CPU round-trip tests, including one using the real openpi-client
  packer (guarded by pytest.importorskip).

Fixes vllm-project#4506. Regression from vllm-project#4282; restores cross-compat from vllm-project#3673.

Signed-off-by: Hyoseop Song <crad_on25@naver.com>
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Nughm3 pushed a commit to Nughm3/vllm-omni that referenced this pull request Jun 18, 2026
tzhouam added a commit to tzhouam/vllm-omni that referenced this pull request Jun 20, 2026
The BDE runner override read `getattr(od_config, "diffusion_model_runner_cls",
None) or platform_hook()`. With a Mock od_config (upstream worker tests from vllm-project#4282,
e.g. test_model_runner_resolved_via_platform / the cuda-profiler tests), the getattr
returns a truthy auto-Mock, shadowing the platform hook so it's never called. Only
honor the override when it's an explicit import-path string, so the platform default
path is unaffected.

Fixes the buildkite CPU "Simple" step (5 tests/diffusion failures); tests/bde +
worker tests green.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Signed-off-by: Taichang Zhou <tzhouam@connect.ust.hk>
tzhouam added a commit to tzhouam/vllm-omni that referenced this pull request Jul 4, 2026
The BDE runner override read `getattr(od_config, "diffusion_model_runner_cls",
None) or platform_hook()`. With a Mock od_config (upstream worker tests from vllm-project#4282,
e.g. test_model_runner_resolved_via_platform / the cuda-profiler tests), the getattr
returns a truthy auto-Mock, shadowing the platform hook so it's never called. Only
honor the override when it's an explicit import-path string, so the platform default
path is unaffected.

Fixes the buildkite CPU "Simple" step (5 tests/diffusion failures); tests/bde +
worker tests green.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Signed-off-by: Taichang Zhou <tzhouam@connect.ust.hk>
khairulkabir1661 pushed a commit to khairulkabir1661/vllm-omni that referenced this pull request Sep 25, 2026
Signed-off-by: yuzhud <yuzhudong@gmail.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

ready label to trigger buildkite CI

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants