Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
415 changes: 202 additions & 213 deletions docs/design/fullduplex-personaplex.md

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions docs/models/supported_models.md
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ th {
| `VoxtralTTSForConditionalGeneration` | Voxtral TTS | `mistralai/Voxtral-4B-TTS-2603` | ✅︎ | ✅︎ | | | — |
| `CovoAudioForConditionalGeneration` | Covo-Audio-Chat | `tencent/Covo-Audio-Chat` | ✅︎ | | | | — |
| `MiniCPMO45OmniForConditionalGeneration` | MiniCPM-o 4.5 | `openbmb/MiniCPM-o-4_5` | ✅︎ | | ✅︎ | | [Repository](https://github.com/vllm-project/vllm-omni/blob/main/recipes/OpenBMB/MiniCPM-o-4_5.md) |
| `PersonaPlexTalkerForConditionalGeneration` | PersonaPlex (full duplex, 24 kHz speech in/out over `/v1/realtime?duplex=1`) | `nvidia/personaplex-7b-v1` | ✅︎ | | | | [Repository](https://github.com/vllm-project/vllm-omni/blob/main/recipes/NVIDIA/PersonaPlex.md) |
| `ErnieImagePipeline` | ERNIE-Image | `baidu/ERNIE-Image`, `baidu/ERNIE-Image-Turbo` | ✅︎ | ✅︎ | ✅︎ | ✅︎ | — |
| `GepardTalkerForConditionalGeneration` | Gepard-1.0 | `nineninesix/gepard-1.0` | ✅︎ | | | ✅︎ | — |
| `HiDreamImagePipeline` | HiDream-I1-Full | `HiDream-ai/HiDream-I1-Full` | ✅︎ | ✅︎ | | | — |
Expand Down
12 changes: 8 additions & 4 deletions docs/serving/realtime_duplex_api.md
Original file line number Diff line number Diff line change
Expand Up @@ -384,9 +384,9 @@ with no OpenAI counterpart.

The event vocabulary is uniform, but several surfaces are gated by the
`capabilities` object the server returns in `session.created`; a client must
branch on those flags rather than on the model name. MiniCPM-o 4.5 is the
only model on the plugin contract today; the other two columns record what
their integrations advertise once the follow-up PRs port them:
branch on those flags rather than on the model name. MiniCPM-o 4.5 and
PersonaPlex are on the plugin contract today; the Nemotron VoiceChat column
records what its integration advertises once the follow-up PR ports it:

| Capability | MiniCPM-o 4.5 | PersonaPlex | Nemotron VoiceChat | Gated surface |
| --- | --- | --- | --- | --- |
Expand All @@ -401,7 +401,11 @@ their integrations advertise once the follow-up PRs port them:
Everything else in the catalogue — session lifecycle, heartbeat and event
acknowledgement, append/commit/clear, the response envelope, playback
acknowledgement, and the error envelope — behaves identically for every
model.
model. Two PersonaPlex specifics follow from its capabilities rather than from
special-casing: a model with `supports_client_commit=false` auto-responds
without `extra_body.auto_response`, and `response.cancel` /
`output_audio_buffer.clear` restart its conversation context (a new Stage 0
request replays the voice/persona prefill).

### Compatibility with the OpenAI Realtime protocol

Expand Down
56 changes: 36 additions & 20 deletions examples/online_serving/personaplex/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,15 +2,15 @@

Serve [`nvidia/personaplex-7b-v1`](https://huggingface.co/nvidia/personaplex-7b-v1)
(a Moshi-based full-duplex speech-to-speech model) with the native vLLM-Omni engine
through the unified duplex serving stack (`/v1/duplex` and `/v1/realtime?duplex=1`).
through the unified full-duplex framework (`/v1/realtime?duplex=1`, alias `/v1/duplex`).

> Requires a GPU and Hugging Face access to the gated repo
> (`HF_TOKEN` with access to `nvidia/personaplex-7b-v1`).

## Start the server

The default `vllm_omni/deploy/personaplex.yaml` enables the engine-owned
full-duplex control plane (`session_mode: duplex`):
The default `vllm_omni/deploy/personaplex.yaml` is a duplex deployment
(`session_mode: duplex`, two sessions per replica):

```bash
HF_TOKEN=... CUDA_VISIBLE_DEVICES=0 python -m vllm_omni.entrypoints.cli.main serve \
Expand All @@ -19,18 +19,32 @@ HF_TOKEN=... CUDA_VISIBLE_DEVICES=0 python -m vllm_omni.entrypoints.cli.main ser
--deploy-config vllm_omni/deploy/personaplex.yaml
```

This exposes:
This exposes `WS /v1/realtime?duplex=1` (alias `WS /v1/duplex`): the OpenAI
Realtime session protocol projected onto vLLM-Omni duplex sessions (client API and
wire protocol: [`docs/serving/realtime_duplex_api.md`](../../../docs/serving/realtime_duplex_api.md)).
There is no `/v1/chat/completions` route: PersonaPlex answers speech only.

- `WS /v1/duplex` — the native duplex session dialect
(`session.create` / `input_audio_buffer.append` / `response.output_audio.delta` ...);
- `WS /v1/realtime?duplex=1` — the same sessions projected onto the
OpenAI Realtime event vocabulary (client API and wire protocol:
[`docs/serving/realtime_duplex_api.md`](../../../docs/serving/realtime_duplex_api.md)).
PersonaPlex is a pure-lockstep model: every session is native duplex, audio flows
continuously in both directions in 80 ms frames, the model decides when to speak,
and there are no client commits or external turn signals
(`supports_client_commit=false`, `supports_external_turn_signal=false`). A session
therefore auto-responds without any vendor flag.

PersonaPlex is a pure-lockstep model: every session on a PersonaPlex deployment is
native duplex (`is_enabled()` is unconditionally true), audio flows continuously in
both directions, and there are no client commits or external turn signals
(`supports_client_commit=false`, `supports_external_turn_signal=false`).
## Talk to it

With the client library and the PersonaPlex preset (24 kHz `pcm_f32le` in, bundled
voice prompt, persona text):

```python
from vllm_omni.clients.duplex import DuplexClient
from vllm_omni.clients.personaplex import create_duplex_session_config

cfg = create_duplex_session_config(voice="NATF2.pt", persona="You are a concise assistant.")
async with DuplexClient("ws://127.0.0.1:8000/v1/realtime?duplex=1", model="/path/to/personaplex-7b-v1", config=cfg) as c:
await c.stream_pcm(pcm_f32le_24k) # keep streaming; the model speaks while it listens
```

Voice and persona are fixed for the session (`session.update` cannot change them).

## Validate the serving path

Expand All @@ -45,9 +59,11 @@ python tests/e2e/online_serving/personaplex_realtime_duplex.py \
--output-dir /tmp/personaplex-realtime-duplex
```

The unified endpoint advertises `supports_barge_in=false`: overlapping speech
is native model behavior, but destructive output interruption and model-state
rewind have not been validated for PersonaPlex.
The endpoint advertises `supports_barge_in=false`: overlapping speech is native
model behaviour, but destructive output interruption and model-state rewind have
not been validated for PersonaPlex. `response.cancel` and
`output_audio_buffer.clear` restart the model's conversation context (a fresh
Stage 0 request replays the voice/persona prefill).

## Notes

Expand All @@ -56,8 +72,8 @@ rewind have not been validated for PersonaPlex.
regardless of engine speed. On localhost it is smooth.
- The earlier standalone Moshi-web compatibility server (browser client at `/`,
binary WS protocol at `/api/chat`, raw-PCM `/v1/audio/duplex`) was demo-only
and has been removed; use the unified endpoints above.
- Session config (voice / persona / sampling) is passed per session via
`extra_body`; see
`vllm_omni/model_executor/models/personaplex/duplex/serving_adapter.py`.
and has been removed; use the unified endpoint above.
- The model plugin, worker-side lockstep runtime and input framing live in
`vllm_omni/model_executor/models/personaplex/duplex/`; design notes in
[`docs/design/fullduplex-personaplex.md`](../../../docs/design/fullduplex-personaplex.md).
- Full runbook: [`recipes/NVIDIA/PersonaPlex.md`](../../../recipes/NVIDIA/PersonaPlex.md).
22 changes: 12 additions & 10 deletions recipes/NVIDIA/PersonaPlex.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,10 +25,10 @@ matching the reference implementation frame for frame on the golden replays used
as the acceptance gate.

PersonaPlex is the first Moshi-class (pure-lockstep) model on the vLLM-Omni
full-duplex serving stack: it plugs into the generic `/v1/duplex` handler
through the standard plugin seams (`duplex_serving_adapter` /
`duplex_runtime_extension` in its `pipeline.py`) with a model-specific
package at `vllm_omni/model_executor/models/personaplex/duplex/`.
unified full-duplex framework: its `pipeline.py` declares one
`duplex_plugin` (`PersonaPlexDuplexPlugin`), and the model-specific package at
`vllm_omni/model_executor/models/personaplex/duplex/` holds the plugin, the
worker-side lockstep Stage 0 runtime and the 80 ms input framing.

## References

Expand Down Expand Up @@ -114,16 +114,18 @@ HF_TOKEN=... CUDA_VISIBLE_DEVICES=0 python -m vllm_omni.entrypoints.cli.main ser
--deploy-config vllm_omni/deploy/personaplex.yaml
```

This exposes `WS /v1/duplex` (native duplex dialect) and
`WS /v1/realtime?duplex=1` (OpenAI Realtime projection; client API and wire
protocol in [`docs/serving/realtime_duplex_api.md`](../../docs/serving/realtime_duplex_api.md)).
Voice and persona are set per session via `extra_body`.
This exposes `WS /v1/realtime?duplex=1` (alias `WS /v1/duplex`): the OpenAI
Realtime session protocol, client API and wire vocabulary in
[`docs/serving/realtime_duplex_api.md`](../../docs/serving/realtime_duplex_api.md).
Voice (`voice`, a bundled `.pt` basename) and persona (`instructions`) are set
in `session.update`; the model takes no client commits and serves no
`/v1/chat/completions` route.

#### Verification

```bash
# GPU-free contract tests (stage0 runtime + unified serving adapter)
pytest tests/model_executor/models/personaplex/duplex/ -q
# GPU-free contract tests (plugin, stage0 runtime, runner scenario)
pytest tests/model_executor/models/personaplex/duplex/ tests/engine/duplex/test_session_runner_personaplex.py -q

# GPU e2e: paced 24 kHz PCM over /v1/realtime?duplex=1, two concurrent
# sessions, overflow admission, slot recycling, non-silent output
Expand Down
64 changes: 33 additions & 31 deletions tests/e2e/online_serving/personaplex_realtime_duplex.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
import json
import math
import time
import uuid
import wave
from collections.abc import Awaitable, Callable, Sequence
from pathlib import Path
Expand Down Expand Up @@ -125,10 +124,11 @@ def _input_identity(
return {"path": str(resolved), "sha256": actual}


def _realtime_url(base_url: str, model: str, session_id: str) -> str:
def _realtime_url(base_url: str, model: str) -> str:
"""The duplex Realtime URL; the session id is allocated by the server, never chosen here."""
parts = urlsplit(base_url)
query = dict(parse_qsl(parts.query, keep_blank_values=True))
query.update(duplex="1", model=model, autostart="0", session_id=session_id)
query.update(duplex="1", model=model, autostart="0")
return urlunsplit((parts.scheme, parts.netloc, parts.path, urlencode(query), parts.fragment))


Expand Down Expand Up @@ -197,16 +197,27 @@ def _audio_bytes(client: RawRealtimeProbe) -> bytes:
return b"".join(chunk for _, chunk in _validated_audio_chunks(client) if chunk)


def _session_id(created: dict[str, object]) -> str | None:
"""The server-allocated id announced in ``session.created``."""
session = created.get("session")
if isinstance(session, dict):
for key in ("id", "session_id"):
value = session.get(key)
if isinstance(value, str) and value:
return value
value = created.get("session_id")
return value if isinstance(value, str) and value else None


async def _open_session(
args: argparse.Namespace,
*,
session_id: str,
persona: str,
expect_error: bool = False,
) -> tuple[RawRealtimeProbe, dict[str, object]]:
client = RawRealtimeProbe(_realtime_url(args.url, args.model, session_id))
client = RawRealtimeProbe(_realtime_url(args.url, args.model))
await client.__aenter__()
await client.send(_session_update(args, session_id=session_id, persona=persona))
await client.send(_session_update(args, persona=persona))
event_type = "error" if expect_error else "session.created"
await wait_for(
lambda: client.events.count(event_type) > 0,
Expand All @@ -216,11 +227,10 @@ async def _open_session(
return client, _events(client, event_type)[-1]


def _session_update(args: argparse.Namespace, *, session_id: str, persona: str) -> dict[str, object]:
def _session_update(args: argparse.Namespace, *, persona: str) -> dict[str, object]:
return {
"type": "session.update",
"session": {
"session_id": session_id,
"model": args.model,
"modalities": ["audio", "text"],
"input_audio_format": "pcm_f32le",
Expand Down Expand Up @@ -426,13 +436,11 @@ async def run(args: argparse.Namespace) -> dict[str, object]:
output_dir = Path(args.output_dir)
output_dir.mkdir(parents=True, exist_ok=True)

ids = {name: f"personaplex-{name}-{uuid.uuid4().hex}" for name in ("primary", "secondary", "replacement")}
primary, created = await _open_session(args, session_id=ids["primary"], persona=args.persona)
secondary, secondary_created = await _open_session(
args,
session_id=ids["secondary"],
persona=args.secondary_persona,
)
primary, created = await _open_session(args, persona=args.persona)
secondary, secondary_created = await _open_session(args, persona=args.secondary_persona)
ids = {"primary": _session_id(created), "secondary": _session_id(secondary_created)}
if not ids["primary"] or not ids["secondary"] or ids["primary"] == ids["secondary"]:
raise AssertionError(f"server did not allocate distinct session ids: {ids}")
capabilities = _capabilities(created)
if _capabilities(secondary_created) != capabilities:
raise AssertionError("concurrent sessions returned different capabilities")
Expand All @@ -446,12 +454,7 @@ async def run(args: argparse.Namespace) -> dict[str, object]:
if any(capabilities.get(key) != value for key, value in expected_capabilities.items()):
raise AssertionError(f"unexpected PersonaPlex capabilities: {capabilities}")

overflow, error = await _open_session(
args,
session_id=f"personaplex-overflow-{uuid.uuid4().hex}",
persona=args.persona,
expect_error=True,
)
overflow, error = await _open_session(args, persona=args.persona, expect_error=True)
error_body = error.get("error")
overflow_code = error_body.get("code") if isinstance(error_body, dict) else error.get("code")
await overflow.__aexit__(None, None, None)
Expand Down Expand Up @@ -492,11 +495,10 @@ async def run(args: argparse.Namespace) -> dict[str, object]:
_save(output_dir, "primary", primary, primary_audio)
await _close_session(primary, timeout_s=args.timeout_s)

replacement, replacement_created = await _open_session(
args,
session_id=ids["replacement"],
persona=args.replacement_persona,
)
replacement, replacement_created = await _open_session(args, persona=args.replacement_persona)
ids["replacement"] = _session_id(replacement_created)
if not ids["replacement"] or ids["replacement"] in {ids["primary"], ids["secondary"]}:
raise AssertionError(f"replacement session id was not freshly allocated: {ids}")
if _capabilities(replacement_created) != capabilities:
raise AssertionError("replacement session returned different capabilities")
continuation_frames, replacement_frames = await asyncio.gather(
Expand Down Expand Up @@ -694,8 +696,9 @@ async def _run_load_session(
ready: asyncio.Future[None],
start: asyncio.Future[float],
) -> dict[str, object]:
session_id = f"personaplex-load-{index}-{uuid.uuid4().hex}"
client = RawRealtimeProbe(_realtime_url(args.url, args.model, session_id), close_timeout_s=args.cleanup_timeout_s)
# Label only: the server allocates the real session id (read back from session.created).
session_id = f"personaplex-load-{index}"
client = RawRealtimeProbe(_realtime_url(args.url, args.model), close_timeout_s=args.cleanup_timeout_s)
sends: list[tuple[float, float, float]] = []
error: str | None = None
cleanup_error: str | None = None
Expand All @@ -706,10 +709,9 @@ async def _run_load_session(
try:
async with client:
try:
await asyncio.wait_for(
client.send(_session_update(args, session_id=session_id, persona=args.persona)), args.timeout_s
)
await asyncio.wait_for(client.send(_session_update(args, persona=args.persona)), args.timeout_s)
await _wait_load_event(client, "session.created", args.timeout_s)
session_id = _session_id(_events(client, "session.created")[-1]) or session_id
capabilities = _capabilities(_events(client, "session.created")[-1])
if (
capabilities.get("chunk_period_ms") != 80
Expand Down
84 changes: 84 additions & 0 deletions tests/e2e/online_serving/test_personaplex_duplex.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project

"""GPU coverage for PersonaPlex on the unified full-duplex framework.

Boots the PersonaPlex deploy (``vllm_omni/deploy/personaplex.yaml``) and runs
the strict Realtime driver (``personaplex_realtime_duplex.py``): two paced
24 kHz sessions, admission overflow, per-session slot recycling, audible
whole-frame output. The checkpoint is gated, so the test runs only when
``PERSONAPLEX_MODEL_PATH`` points at a local copy of ``nvidia/personaplex-7b-v1``.
"""

from __future__ import annotations

import asyncio
import os
from pathlib import Path

import pytest

from tests.e2e.online_serving import personaplex_realtime_duplex as driver
from tests.helpers.mark import hardware_test
from tests.helpers.runtime import OmniServerParams
from tests.helpers.stage_config import get_deploy_config_path

pytestmark = pytest.mark.omni

MODEL_PATH = os.environ.get("PERSONAPLEX_MODEL_PATH", "")
DEPLOY_CONFIG = get_deploy_config_path("personaplex.yaml")

SERVER_PARAMS = [
pytest.param(
OmniServerParams(
model=MODEL_PATH or "nvidia/personaplex-7b-v1",
stage_config_path=DEPLOY_CONFIG,
use_stage_cli=False,
server_args=["--trust-remote-code"],
),
id="two-stage-single-gpu",
)
]

requires_checkpoint = pytest.mark.skipif(
not MODEL_PATH or not Path(MODEL_PATH).is_dir(),
reason="set PERSONAPLEX_MODEL_PATH to a local nvidia/personaplex-7b-v1 checkout",
)


def _speech_wav() -> Path:
"""Real speech (the MiniCPM-o asset, 16 kHz; the driver resamples to 24 kHz).

PersonaPlex answers speech, not tones: a synthetic signal yields a silent
reply and the driver's audibility floor rightly fails it.
"""
return Path(__file__).resolve().parents[2] / "assets" / "minicpmo_4_5" / "response_required_16k.wav"


@requires_checkpoint
@pytest.mark.advanced_model
@hardware_test(res={"cuda": "H100"}, num_cards=1)
@pytest.mark.parametrize("omni_server", SERVER_PARAMS, indirect=True)
def test_personaplex_realtime_duplex_sessions(omni_server, tmp_path: Path) -> None:
input_wav = _speech_wav()
args = driver.parse_args(
[
"--url",
f"ws://{omni_server.host}:{omni_server.port}/v1/realtime?duplex=1",
"--model",
MODEL_PATH,
"--input-wav",
str(input_wav),
"--output-dir",
str(tmp_path / "out"),
]
)

result = asyncio.run(driver.run(args))

sessions = {name: result[name] for name in ("primary", "secondary", "replacement")}
assert all(isinstance(session, dict) for session in sessions.values())
ids = {name: session["session_id"] for name, session in sessions.items() if isinstance(session, dict)}
assert ids["primary"]
assert ids["secondary"] != ids["primary"]
assert ids["replacement"] not in {ids["primary"], ids["secondary"]}
Loading
Loading