Repository navigation
[Feature] JoyAI-VL-Interaction streaming interaction serving layer - #4575
Conversation
|
Codex usage limits have been reached for code reviews. Please check with the admins of this repo to increase the limits by adding credits. |
e17bc73 to
fd77975
Compare
f62f13b to
93df450
Compare
Add vllm_omni/experimental/fullduplex/: a vLLM-Omni-native serving layer for proactive streaming video-language interaction (per-second speak / silence / delegate). Includes a generic full-duplex core (DuplexRuntime/Session/Adapter) and a JoyVL implementation: per-tick decision policy, 3-tier summary memory, OpenAI-compatible orchestrator, and pluggable ASR / TTS / delegation bridges (chat, text-to-image, image-edit, and a routing bridge). Ships the JD recipe, demo glue under examples/, and unit tests. Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com>
93df450 to
ba01d34
Compare
linyueqian
left a comment
There was a problem hiding this comment.
Reviewed end to end. I reconstructed the package and ran it (this is the public version of the framework I looked at earlier).
Verified
tests/fullduplex21/21 pass under the newexperimental/layout.- Serving stack imports and the FastAPI app builds. I drove
InteractionSession.stepboth against a fake backend and against the realJoyAI-VL-Interaction-Previewmodel: forced-silence-before-query (skips the model call), response dedup -> silence, delegate submit+fold, async 3-tier memory, andskip_special_tokens=Falseall behave correctly. End to end the orchestrator reproduces the speak/silence decisions faithfully. app.pyand thertsp/scripts are clean and talk to the orchestrator over HTTP only.
Good calls
- Moving under
experimental/, and the newREADME.md"Scope" note thatcore/+adapter.pyare a demonstration (exercised by the tests), not the path the HTTP orchestrator currently drives. That clears up the main thing I would have flagged: the serving path isjoyvl/serving/->decision/+memory/directly, whileDuplexRuntime/JoyVLDuplexAdapterare referenced only intest_runtime.py. Documenting it plainly is the right call. SessionManager.reset()now callssession.reset(), so in-flight consolidation tasks are cancelled on reset/evict.- RTSP ingestion is a nice way around browser/WebRTC transport limits.
Minor
- [low]
_completion_responseemits a non-standard top-levelstreamingharnesskey alongsideinteraction; it duplicatesinteraction.memory+ timing. Consider renaming or dropping it. - [low] The PR description is a single line. For a 3.4k-line feature a short summary plus the reproduce steps (they already exist in the recipe) in the body would help reviewers; the PR body is the first thing people read.
- [low, env-dependent] Deploy note worth adding to the recipe: on a host without nvcc/ninja, a plain
vllm serveof the 8B can crash engine-core in the FlashInfer sampler JIT (FileNotFoundError: 'ninja') duringprofile_run.VLLM_USE_FLASHINFER_SAMPLER=0(or installing ninja) works around it.
Overall this is in good shape and the logic is correct against the model. Nothing blocking from me.
Gaohan123
left a comment
There was a problem hiding this comment.
Great work! Let's unify the api-server and full-duplex runtime in the following PRs
| @@ -0,0 +1,131 @@ | |||
| # SPDX-License-Identifier: Apache-2.0 | |||
| # SPDX-FileCopyrightText: Copyright contributors to the vLLM project | |||
|
|
|||
There was a problem hiding this comment.
I think later we can integrate asr and tts into stage_client of vllm-omni to better utillize features like async_chunk in follow-up PRs. Like PR #4257
There was a problem hiding this comment.
Will keep in mind.
| --limit-mm-per-prompt '{"image":256,"video":1}' | ||
|
|
||
| # 2. Interaction orchestrator (OpenAI-compatible, :8070) | ||
| python -m vllm_omni.experimental.fullduplex.joyvl.serving.server --port 8070 \ |
There was a problem hiding this comment.
The transitional solution is ok. I suggest we need to integrate fulllduplex server into vllm-omni api-server in following PRs
There was a problem hiding this comment.
Agreed — this is a transitional standalone orchestrator; folding the full-duplex server into the vLLM-Omni api-server is the right next step and I'll take it in a follow-up PR.
| orchestrator over its HTTP API only: | ||
|
|
||
| ```bash | ||
| pip install vllm-omni[demo] # gradio + opencv-python + requests |
There was a problem hiding this comment.
Better use uv with consistency
|
|
||
| from vllm_omni.experimental.fullduplex.joyvl.decision.output_parser import parse_action | ||
| from vllm_omni.experimental.fullduplex.joyvl.memory.brain import InteractionBrain | ||
|
|
There was a problem hiding this comment.
The pytest marks is incorrect. Please refer to https://docs.vllm.ai/projects/vllm-omni/en/latest/contributing/ci/tests_markers/
| OpenAIDelegationBridge, | ||
| ) | ||
| from vllm_omni.experimental.fullduplex.joyvl.decision.policy import JoyVLPolicy, sample_frames | ||
|
|
There was a problem hiding this comment.
Fixed in 6984e11 — same core_model / cpu pytestmark added here.
| from vllm_omni.experimental.fullduplex.core.runtime import DuplexRuntime | ||
| from vllm_omni.experimental.fullduplex.core.session import DuplexSession, DuplexSessionConfig | ||
| from vllm_omni.experimental.fullduplex.joyvl.adapter import JoyVLDuplexAdapter | ||
|
|
|
|
||
| ## Scope | ||
|
|
||
| The runnable serving path is `joyvl/serving/` driving `decision/` + `memory/` directly. |
There was a problem hiding this comment.
Important: left for following PRs
…io app Reference the upstream JoyAI-VL-Interaction WebUI (services/webui) in front of the orchestrator instead of shipping a separate Gradio app. Remove app.py and gradio from the demo extra. Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com>
a12a8e8 to
5ef25db
Compare
hsliuustc0106
left a comment
There was a problem hiding this comment.
Posted the inline review notes we discussed.
|
|
||
| async def _start_response(self, prev: asyncio.Task | None, emit: Emit) -> asyncio.Task: | ||
| if prev is not None and not prev.done(): | ||
| await prev |
There was a problem hiding this comment.
[blocking] _start_response() waits for the previous response to finish before the input loop can continue. That means a second response trigger can delay later response.cancel / barge-in events until the old response has fully streamed, which breaks the full-duplex cancellation contract. I reproduced this with a slow adapter: the first response emitted all deltas before the cancel was handled. This should cancel/epoch-stale the active response, or otherwise keep the input loop able to process cancel events immediately.
There was a problem hiding this comment.
Fixed in 7abcc07. _start_response no longer awaits the previous response's full stream — it stales the active response (barge-in) and cancels it, then starts the new one, so the input loop can service cancel/barge-in events immediately. Added test_runtime_new_response_supersedes_inflight_without_blocking (slow adapter, second trigger supersedes the first; only the latest response completes).
| ) | ||
| if url: | ||
| return chat_bridge(url) | ||
| return StubDelegationBridge() |
There was a problem hiding this comment.
[blocking] With the current defaults, omitting --delegation-backend-url still enables delegation through StubDelegationBridge(). That contradicts the recipe guidance that leaving the URL unset keeps delegation off, and in real use it can fold fake stub answers into session memory. The default should probably be no delegation unless a real backend is configured, with the stub limited to tests/demo-only paths.
There was a problem hiding this comment.
Fixed in 7abcc07. _build_delegation now returns None (delegation off) unless a real backend is configured; StubDelegationBridge is opt-in only via --delegation-kind stub. So omitting --delegation-backend-url truly keeps delegation off — no stub answers folded into memory — matching the recipe. Added test_build_delegation_off_unless_backend_configured.
| async def submit(self, question: str, note: str, frames: list[tuple[str, str]]) -> str: | ||
| self._counter += 1 | ||
| task_id = f"deleg-{self._counter}" | ||
| self._tasks[task_id] = asyncio.create_task(self._solve(question, note, frames)) |
There was a problem hiding this comment.
[non-blocking] These background tasks are removed only when poll() is called after completion. If a session is reset/evicted or the client disconnects before polling again, the shared bridge keeps the finished task/result around indefinitely. Consider a done callback or explicit reset/cancel path so long-lived demo servers do not accumulate abandoned delegation tasks.
There was a problem hiding this comment.
Fixed in 7abcc07. Added cancel(task_id) to every bridge (and the protocol); brain.reset() now cancels all pending delegations, so reset/evicted sessions don't leave finished-but-unpolled tasks on the shared bridge. Added test_delegation_cancel_drops_task.
|
|
||
|
|
||
| def create_app(config: InteractionConfig) -> FastAPI: | ||
| app = FastAPI(title="vLLM-Omni Interaction Server") |
There was a problem hiding this comment.
No shutdown/lifespan hook, so the AsyncOpenAI/httpx clients and pending background tasks leak on shutdown (and on every uvicorn reload).
create_app only registers routes and main() just runs uvicorn. OpenAIBackend._client and the AsyncOpenAI held by OpenAIDelegationBridge are never closed, and the asyncio.create_tasks in delegation._tasks and session._consolidating are never cancelled — so on graceful stop you get "Unclosed client session" warnings and orphaned tasks.
Suggestion: add a lifespan to FastAPI(...) (or @app.on_event("shutdown")) that awaits an aclose() on the manager, and give OpenAIBackend / the delegation bridges an async def aclose() that closes the client and cancels + awaits the pending tasks.
hsliuustc0106
left a comment
There was a problem hiding this comment.
Adding two follow-up lifecycle/race findings after another pass.
|
|
||
| def create_app(config: InteractionConfig) -> FastAPI: | ||
| app = FastAPI(title="vLLM-Omni Interaction Server") | ||
| manager = SessionManager(config) |
There was a problem hiding this comment.
[should-fix] This app never registers a shutdown/lifespan hook, so the shared async resources owned by SessionManager are left for process teardown: OpenAIBackend._client, OpenAIDelegationBridge._client, delegation _tasks, and each session's _consolidating tasks. On uvicorn shutdown/reload this can leave unclosed httpx clients and pending tasks. Please add a manager aclose() and call it from FastAPI lifespan/shutdown; the backends/bridges should close their clients and cancel+await pending tasks.
There was a problem hiding this comment.
Fixed in 1ee7489. create_app now registers a FastAPI lifespan that calls SessionManager.aclose() on shutdown/reload: it cancels+awaits pending delegation and consolidation tasks and closes the AsyncOpenAI/httpx clients. Added aclose() to ModelBackend/OpenAIBackend, Summarizer, and every DelegationBridge (+protocol). Test: test_manager_aclose_closes_without_error (idempotent).
| async with self._locks[session_id]: | ||
| return await session.step(frames, query) | ||
|
|
||
| def reset(self, session_id: str) -> None: |
There was a problem hiding this comment.
[should-fix] reset() mutates and drops the session without acquiring the per-session lock. A concurrent /v1/chat/completions can be inside session.step() under that lock while /reset calls session.reset() and removes the lock/session from the manager, so the in-flight request continues on a reset/discarded session. Making reset async and acquiring the same lock before resetting/removing the session would make reset semantics deterministic. InteractionSession.reset() should also await cancelled consolidation tasks instead of clearing the set immediately.
There was a problem hiding this comment.
Fixed in 1ee7489. SessionManager.reset/set_persona are now async and acquire the per-session lock, so they wait for any in-flight step() on that session instead of mutating/dropping it mid-request; step captures session+lock together (no await between) so a concurrent reset can't swap them out. InteractionSession.reset()/aclose() now await asyncio.gather(*consolidating, return_exceptions=True) after cancelling, instead of clearing the set immediately. Test: test_session_reset_cancels_and_awaits_consolidation.
| async with self._locks[session_id]: | ||
| return await session.step(frames, query) | ||
|
|
||
| def reset(self, session_id: str) -> None: |
There was a problem hiding this comment.
reset() / set_persona() mutate session state without the per-session lock, so they race an in-flight step() on the same session.
manager.reset() pops _sessions / _locks and calls session.reset() while another coroutine may be awaiting session.step() under that lock. session.reset() also does task.cancel() on the consolidation tasks without awaiting, so they keep running (against a discarded brain) until the next loop tick. Impact is bounded — the session is dropped — but the in-flight request continues on a reset session and its lock is already gone from the dict.
Suggestion: make manager.reset / set_persona acquire self._locks[session_id] (i.e. make them async), and in session.reset() do await asyncio.gather(*self._consolidating, return_exceptions=True) after cancelling.
| ) | ||
| if url: | ||
| return chat_bridge(url) | ||
| return StubDelegationBridge() |
There was a problem hiding this comment.
Defaulting delegation on with the stub contradicts the recipe and folds fake text into memory.
enable_delegation defaults to True and --no-delegation is off by default, so a deployment that omits --delegation-backend-url lands here: StubDelegationBridge. That stub's poll() returns a ready digest of "(background result for: …) — stub digest; wire a real brain here.", which fold_delegations then appends to qa_history and feeds into the memory prefix for every later turn. But the recipe tells users: "Omit --delegation-backend-url to keep delegation off."
Suggestion: when no backend URL is configured, keep delegation off (return None / don't build a bridge, equivalent to --no-delegation), and reserve StubDelegationBridge for tests/demo.
| async def submit(self, question: str, note: str, frames: list[tuple[str, str]]) -> str: | ||
| self._counter += 1 | ||
| task_id = f"deleg-{self._counter}" | ||
| self._tasks[task_id] = asyncio.create_task(self._solve(question, note, frames)) |
There was a problem hiding this comment.
Submitted delegation tasks live in self._tasks until poll() removes them, so abandoned tasks accumulate for the life of the process.
The bridge is a singleton on the SessionManager, shared across all sessions. poll() only runs while a session keeps calling fold_delegations; if a session is reset/evicted (which clears brain._pending_delegations but not this dict) or a client disconnects before folding, the entry — and the still-running create_task — stays here forever.
Suggestion: add a done_callback that drops the entry once the task finishes, and evict by id on reset; or track tasks per-session and cancel+evict on reset.
| await task | ||
| self.session.state = DuplexState.CLOSED | ||
|
|
||
| async def _start_response(self, prev: asyncio.Task | None, emit: Emit) -> asyncio.Task: |
There was a problem hiding this comment.
_start_response does await prev inside the input loop, so a second response trigger blocks all later events — including RESPONSE_CANCEL / barge-in — until the previous response finishes streaming. That breaks the full-duplex contract (you can't interrupt a response): the first slow response drains all its deltas before a cancel is processed.
This is the core/ demonstration path (driven only by test_runtime.py, not the HTTP serving layer), so it's not on the shipped path today — but it's a real bug worth fixing before this runtime informs the api-server merge. Fix: cancel/replace the in-flight response instead of await prev, or move the response stream off the input loop.
…off by default, cancel leaked tasks - runtime: a new response now stales+cancels the in-flight one instead of awaiting its full stream, so the input loop can service cancel/barge-in immediately - serving: _build_delegation returns None unless a real backend is configured; StubDelegationBridge is opt-in via --delegation-kind stub (no more fake answers folded into memory) - delegation: add cancel(task_id) to all bridges; brain.reset() cancels pending delegations so reset/evicted sessions don't leak tasks - tests: supersede-without-blocking, delegation-off-by-default, cancel-drops-task Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com>
|
Thanks @hsliuustc0106 |
…eset/set_persona against in-flight step - create_app now registers a FastAPI lifespan that calls SessionManager.aclose(): closes AsyncOpenAI/httpx clients (backend, summarizer, OpenAIDelegationBridge) and cancels+awaits pending delegation/consolidation tasks on shutdown/reload - SessionManager.reset/set_persona are async and acquire the per-session lock, so they no longer race an in-flight step(); reset removes the session only after the lock is held - InteractionSession.reset()/aclose() cancel AND await consolidation tasks instead of clearing the set immediately - add aclose() to ModelBackend/OpenAIBackend, Summarizer, and all DelegationBridges (+protocol) - tests: delegation aclose, session reset awaits consolidation, manager aclose idempotent Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com>
|
Thanks again @hsliuustc0106 — both follow-up lifecycle findings addressed in 1ee7489: |
|
Note on the |
Signed-off-by: lishunyang <lishunyang12@163.com>
… async_chunk TTS Implements the remaining review follow-ups from vllm-project#4575 (@Gaohan123). B — unify the api-server and full-duplex transport: - joyvl/serving exposes a mountable router (build_interaction_router) and mount_fullduplex(app, config), so the interaction layer runs standalone or is mounted onto the vLLM-Omni OpenAI api-server (opt-in via VLLM_OMNI_FULLDUPLEX_BACKEND_URL). - New WebSocket /v1/fullduplex/stream (DuplexStreamHandler) drives the core DuplexRuntime over the event protocol -> streaming response deltas with barge-in. HTTP (one decision/request) and WS (streaming) share the same InteractionSession engine. C — consume the staged async_chunk TTS: - bridges/speech.py factors the vLLM-Omni staged speech-stream protocol (/v1/audio/speech/stream, qwen3_tts.yaml async_chunk) into a reusable, tested relay; the example tts_bridge uses it and forwards each audio chunk as it arrives. ASR documented as one-shot (staged streaming ASR is the remaining deeper integration). Docs: recipe + fullduplex README updated for the WS transport, the api-server mount, and the staged async_chunk TTS. Tests: tests/fullduplex 34 passed (HTTP route, WS streaming, mount, speech relay). Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com>
…llm-project#4575) Signed-off-by: lishunyang12 <125541396+lishunyang12@users.noreply.github.com> Signed-off-by: Yueqian Lin <linyueqian@outlook.com> Signed-off-by: lishunyang <lishunyang12@163.com> Co-authored-by: Hongsheng Liu <liuhongsheng4@huawei.com> Co-authored-by: Yueqian Lin <linyueqian@outlook.com>
The development of JoyAI-VL 8B was made possible through close collaboration within the JoyAI team
https://arxiv.org/html/2606.14777
https://github.com/jd-opensource/JoyAI-VL-Interaction