Skip to content

[Model] Add single-stage Qwen3-TTS with in-Talker streaming codec decode - #8259

Merged
Sy0307 merged 22 commits into
vllm-project:mainfrom
Sy0307:perf/qwen3-tts-fused-stream
Oct 2, 2026
Merged

Sy0307 merged 22 commits into
vllm-project:mainfrom
Sy0307:perf/qwen3-tts-fused-stream

Conversation

@Sy0307

@Sy0307 Sy0307 commented Sep 28, 2026 •

Copy link
Copy Markdown
Collaborator

Purpose and behavior

Add the opt-in qwen3_tts_fused pipeline for one CUDA GPU. The Talker owns a stateful codec decoder and returns streaming PCM directly, removing the Code2Wav engine and connector hop. The default qwen3_tts pipeline retains two independently placed stages and transport configuration.

  • Run the streaming decoder with Triton, cuBLAS and bounded CUDA graphs on a side stream. Group reference priming by real context length; short references and groups outside captured buckets use eager decoding.
  • Add an explicitly configured fused residual predictor, capacity-aware warmup and shared input validation before dispatch. Existing two-stage deployments with talker_first_audio: true can also use eager-MTP predictor execution.
  • Preserve per-request PCM ownership, cancellation, length termination and first-frame ordering. Preemption saves codec state, MTP RNG, conditioning and exact Talker inputs, rebuilds KV, then restores the codec into the new slot without replaying audio.
  • Pack native two-stage async codec chunks with tensor operations, retaining owned, contiguous, codebook-major CPU-long payloads, legacy lists, reference prefixes and terminal tails. Single-stage PCM bypasses this path.
  • Convert each owned D2H PCM batch to one NumPy view, retain request slices and a pending-sample count, and create Tensor payloads when emitting chunks. Keep BF16/Tensor fallbacks, cancellation, request-ID reuse and direct-first-frame ordering.
  • Reuse assistant text IDs and static CustomVoice prompt pieces, encode PCM16 with bit-exact NumPy operations, and give each API worker its own SO_REUSEPORT listener.

Independent path selection and lifecycle fixes

# One stage on one CUDA GPU.
vllm serve Qwen/Qwen3-TTS-12Hz-1.7B-CustomVoice --omni \
  --deploy-config vllm_omni/deploy/qwen3_tts_fused_single_gpu.yaml

# Two stages on one GPU; this existing experimental profile enables MPS.
vllm serve Qwen/Qwen3-TTS-12Hz-1.7B-CustomVoice --omni \
  --deploy-config vllm_omni/deploy/qwen3_tts_high_concurrency_mrv2_single_gpu.yaml

Single-stage options live in stages[].additional_config; two-stage transport and first-frame options live in connectors.*.extra. Single-stage eligibility has its own module. talker_stream_decode requires a final audio-output Talker, CUDA, MRv2, async chunks, TP=PP=1, an in-process worker and disabled prefix caching; invalid topologies fail explicitly.

The single-stage talker_stream_first_audio queue is opt-in and defaults to false, using PCM from the same streaming decoder. The two-stage talker_first_audio path owns its separate first-frame decoder; the options cannot be enabled together. Deployment documentation includes an overlay that retains predictor settings.

  • Set drop_stale_output at connector-driven prompt replacement; the regression fails if the assignment is removed.
  • Replace the shared scheduler's Qwen architecture check with supports_running_prefix_cache_reset. Fused PCM rejects active running resets before preemption; overrides cannot enable a capability unsupported by the topology.
  • Move reference extraction and grouped priming into model hooks, removing Qwen reference parsing and the fixed 25-frame prime chunk from the shared worker.
  • Set production Base reference context to 72 frames, matching qwen3_tts.yaml. Shorter context remains configurable and can affect timbre continuity. Functional tests cover 25/72; no new full speaker-quality evaluation is claimed.
  • Test fixed-uniform fused-predictor equivalence and nonzero framewise codec output against independent FP32 exact decoding, across the attention window and KV-ring wrap, with eager/graph and FP32/BF16 coverage. The restore regression rejects a no-op restore.

Latest change: revert the new history-write implementation

Base: main@e34f8c4f42f2cf800f1037dfb6014a20dfa25e65. Final head: 6f003b86b30412a7256f06354c2463fc32eafe0f.

This head reverts the slot-owned Triton input-history pool introduced by cf956be, including its kernel and dedicated tests. It restores the existing per-active-request foreach history path and retains the NumPy PCM optimization. Exact-input history remains necessary for preemption replay; it has not been deleted. For 1.7B/max length 4096, allocation is approximately 16 MiB per active request, up to approximately 2 GiB at 128 active requests; budget this with codec/output state after KV profiling. The former unconditional configured-slot reservation is gone.

No experiment artifacts, new Python examples, dependencies or environment switches are added to this PR. The separate step-edge/inline transport research is not part of this implementation.

Current-head validation and performance

Fresh unprofiled serving on physical H200 GPU 3 (GPU-b48a6a9b-33d1-35ae-a2f7-506bbccd9ab0): full Seed-TTS EN1088, C128, CustomVoice 1.7B/Vivian/English, four API workers, production single-stage YAML, capacity 128, FULL graphs, token budget 512/max length 4096. Server CPUs 16-31,112-127, client CPUs 48-63,144-159. Two independent starts, four formal rounds per start; exclude a full priming round and use 64 warmups per formal round. Dataset seed 0, unseeded generation. Software: vLLM 0.30.0, torch 2.13.0+cu130, Transformers 5.14.1, driver 580.173.02/CUDA 13.0. The shared host is not exclusive; retain every low formal round.

Current source Pooled audio-s/s Mean first audio Mean completion
6f003b8, PCM optimization with original exact history 811.17 59.58 ms 719.99 ms

All 8704/8704 formal requests succeeded. Eight rounds: 800.04 / 815.42 / 808.46 / 826.12 / 812.65 / 812.74 / 784.21 / 831.81. Pooled throughput is total audio seconds divided by measured wall time. This is a fresh validation of the reverted implementation, not a matched performance comparison with cf956be; its former 845.76 score does not describe this head.

  • 56 runner tests passed after the revert. Repository-pinned Ruff 0.14.10 check/format and all hooks applicable to the changed Python/YAML files passed, including mypy 3.10, SPDX and forbidden-import checks. Full hook environment initialization encountered unrelated Go/network failures; the scoped configuration excludes only actionlint/Markdown hooks that match none of these files. This is not a whole-repository hook-pass claim.
  • All 4296 tracked source-file hashes match the actual serving snapshot. Contributor author/committer identity and DCO sign-offs pass; whitespace checks pass.
  • A real-weight forced-KV-preemption audit passes 33 HTTP requests, 16 suspend/resume/exact-replay events, with history/suspended/restoring/position state empty afterward. It uses capacity 8, max length 512 and 64 KV blocks to force the path; it does not estimate normal production preemption frequency.
  • Six fresh cancellation-diagnosis rounds each on this implementation and parent bd950f7, 456 HTTP requests total, pass every sequential streaming/WAV pair, immediate recovery and recovery after a 1-second drain. Raw per-request PCM is retained locally.
  • Two immediate-after-cancel PCM equality checks failed during the main current-source performance matrix. The original validator did not retain the failed recovery audio, so the cause remains unclassified. Later diagnosis did not reproduce the failures; this does not establish a fix. Fixed seeds do not guarantee batch-layout-invariant PCM. The validator was extended to retain audio without relaxing its assertions. Earlier optional-queue and old-head observations with discarded PCM also remain unclassified.

No fresh WER/speaker-similarity or matched peak-memory evaluation is claimed. Physical GPUs 4/5 were never used. Raw audio, source manifests, traces, reports and benchmark results remain outside the PR.

Earlier matched component and path evidence

These are historical frozen-source comparisons, not new-head measurements. All use the maintained full EN1088 client, exclude priming and retain low formal rounds.

GPU 1 matched single-stage source C64 pooled audio-s/s / mean first C128 pooled audio-s/s / mean first
bd950f7 604.78 / 33.46 ms 794.09 / 60.14 ms
PCM-only component snapshot 608.75 / 33.03 ms 837.44 / 57.32 ms

Each source had two starts and four formal rounds per concurrency. The retained PCM component showed about +5.46% at C128 in that matched experiment; different GPUs/host conditions prevent combining it with the new 811.17 result.

The earlier GPU 6 path comparison at bd950f7 gave single-stage 795.74 / 59.19 ms versus native two-stage 522.01 / 76.78 ms at C128, with 8704/8704 formal requests successful. Decoder and transport implementations also differ, so the +52.44% throughput result does not isolate stage count. The experimental src25 reproduction reached 870.25 / 48.68 ms, but its pinned-buffer/deferred-output lifecycle differs and is not a production baseline for this PR. Historical approximately 29 ms means were C64 under a different raw-PCM client; they cannot be paired with C128 throughput peaks.

Native two-stage tensor-packing A/B on GPU 6, 0a8e794 versus cbb3a8a, control/packed/packed/control startup order, two formal C128 rounds per start:

Source Four rounds audio-s/s Pooled Mean first
Control 470.53 / 461.61 / 469.82 / 463.25 466.27 67.51 ms
Tensor packing 520.17 / 526.51 / 524.58 / 543.63 528.59 74.42 ms

All 8704 requests succeeded. Packing gives +13.37% throughput and +6.91 ms first-audio latency; it does not improve both metrics. The bit-exact CPU microbenchmark is T25 without reference 510.23→29.14 μs and T1 plus 25 reference frames 67.93→48.79 μs. Single-stage bypasses this packing path, so gains are not additive.

At parent bd950f7, 716 relevant tests passed with one skip; the final Base-72 deployment assertion passed in a separate overlapping 19-test rerun. Earlier serving checks covered default/optional-first single-stage, native two-stage and Base 25/72; these are parent-source evidence. No quality claim is extended beyond the tested source/context.

Reproduce the current benchmark

Select an available physical CUDA GPU before remapping. Use local CustomVoice model weights and the complete Seed-TTS dataset, pin the server/client CPUs as above, start the production YAML with --served-model-name qwen-perf --api-server-count 4 --disable-log-stats --host 127.0.0.1 --port 18913. Exclude one full priming round, then run four formal rounds per startup and repeat at a second startup:

vllm-omni bench serve --omni --host 127.0.0.1 --port 18913 \
  --model qwen-perf --tokenizer /path/to/Qwen3-TTS-12Hz-1.7B-CustomVoice \
  --trust-remote-code --backend openai-audio-speech --endpoint /v1/audio/speech \
  --dataset-name seed-tts-text --dataset-path /path/to/seedtts_testset \
  --seed-tts-locale en --num-prompts 1088 --num-warmups 64 \
  --max-concurrency 128 --request-rate inf --seed 0 \
  --extra-body '{"task_type":"CustomVoice","voice":"Vivian","language":"English","temperature":0.9,"top_p":1.0,"top_k":50,"repetition_penalty":1.05,"max_new_tokens":2048}' \
  --percentile-metrics ttft,e2el,audio_rtf,audio_ttfp,audio_duration,audio_underrun \
  --save-result --save-detailed --result-dir /path/to/results

Regression commands (Linux, matching dependencies; CUDA cases need a CUDA GPU):

cd tests
python -m pytest -q config/test_config_merge.py config/test_omni_config.py \
  core/sched/test_omni_ar_scheduler_stale_drain.py worker_v2/test_omni_ar_model_runner.py \
  -m 'core_model and cpu' --run-level=core_model
python -m pytest -q model_executor/models/qwen3_tts/test_first_audio_graphs.py \
  model_executor/models/qwen3_tts/test_code_predictor_sampling.py \
  -m 'core_model and cuda' --run-level=core_model

Whole-PR inventory and CI

The complete e34f8c4...6f003b8 diff changes 42 files, +3508/-90, net +3418. The remaining profile-driven increment over bd950f7 changes two files, +85/-7 (streaming_audio.py and runner tests). The latest revert changes four files, +3/-216.

Responsibility Files Added / removed
Production deployment YAML 1 +54 / -0
Documentation 1 +64 / -1
Model/audio and stage input 9 +1424 / -37
Shared runtime, API and configuration 18 +881 / -45
Regression tests 13 +1085 / -7

Parent CUDA CI #16504 passed at bd950f7. Old-head CUDA CI #16540 failed; the inspected Qwen3-TTS CustomVoice job failed downloading model weights with HTTP 401 and an expired ci_token, before inference. Other failed jobs have not all been classified.

New-head CUDA CI #16544 was triggered by reapplying ready; Buildkite confirms commit_id=6f003b86b30412a7256f06354c2463fc32eafe0f. GitHub marks the check failed while Buildkite is failing with other jobs still running at this update. The two inspected CustomVoice jobs and the Base job fail on HTTP 401 / expired ci_token before serving becomes ready; the inspected CPU engine and model-executor jobs also show 10 and 12 download-related HTTP-401 failures, with 3643 and 3260 tests passing respectively. This does not classify every failed job or establish CUDA inference success; CI needs a valid model-download credential.

New-head pre-commit, DCO and Python 3.11/3.12 wheel builds passed. AMD/Intel CI failed with causes not classified here; NPU and documentation checks are pending. The PR is ready for review and mergeable, with no approval claimed. Existing threads have not been resolved on reviewers' behalf.

Add the qwen3_tts_fused pipeline: the Talker decodes each frame with a
stateful streaming codec decoder on a side stream and returns PCM as its
final output, without a Code2Wav stage or connector. Each step delivers
the previous step's PCM so the host never waits on the codec; length-ended
requests flush in the same step, and frames of a request missing from the
current batch stay queued. Voice-clone streams are primed with reference
codes.

Also add an opt-in fused residual code predictor, a bit-exact numpy PCM16
encoder, per-API-server SO_REUSEPORT listeners, API-side assistant text
ids, per-stage additional_config model options for single-stage
deployments, and a device drain before the pinned index ring wraps.

Signed-off-by: Sy03 <1370724210@qq.com>
Drop the Code2Wav stateful-decoder option (no profile enabled it), the
unreachable inline stream-decode branch and the old per-request stream
PCM branch in async-chunk output building. Release stream PCM state in
finish_requests instead of a call-count sweep, and replace getattr
defaults with declared attributes so missing state fails loudly.

Signed-off-by: Sy03 <1370724210@qq.com>
@hsliuustc0106 hsliuustc0106 added the tts code related to tts models label Sep 29, 2026
Signed-off-by: Sy03 <1370724210@qq.com>
@Sy0307 Sy0307 added ready label to trigger buildkite CI merge-test label to trigger buildkite merge test CI and removed ready label to trigger buildkite CI labels Sep 29, 2026
@Sy0307
Sy0307 marked this pull request as ready for review September 29, 2026 19:11
@vllm-omni-review-bot

Copy link
Copy Markdown

This PR appears to belong to: docs/design/module/ar_runtime.md.

Module owners: @tzhouam @fake0fan @Gaohan123

Routing: @tzhouam via module of the changed files, CODEOWNERS; @fake0fan via module of the changed files; @Gaohan123 via module of the changed files

@Sy0307, please review your own changes and leave a short self-review comment describing what you checked. PRs without author self-review may not be assigned a reviewer.

Please take a look when you have a chance. If you would like an automated review, mention @vllm-omni-review-bot in a comment.

@vllm-omni-review-bot

vllm-omni-review-bot commented Sep 29, 2026 •

Copy link
Copy Markdown

Omni ReviewBot triage note

Automated triage of commit 43f9bfa0186c produced:

  • Priority: high. Prompt maintainer attention is suggested.

These are automated triage suggestions only — the final decision belongs to the maintainers.

@Sy0307 Sy0307 added ready label to trigger buildkite CI and removed ready label to trigger buildkite CI labels Sep 29, 2026
@hsliuustc0106 hsliuustc0106 added the high priority high priority issue, needs to be done asap label Sep 29, 2026
@Sy0307 Sy0307 added ready label to trigger buildkite CI and removed ready label to trigger buildkite CI labels Sep 30, 2026
@vllm-omni-review-bot

vllm-omni-review-bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Omni ReviewBot: three questions on the performance claim

@Sy0307 this PR reads as a performance or value claim:

  • claim: C128 single-stage throughput is 52.44% higher and mean first audio 17.60 ms lower in this path-level comparison.
  • claim: Packing improves pooled throughput 13.37% and increases mean first audio 6.91 ms; it does not improve both metrics.

Before the full evidence checklist, three short questions:

  1. Bottleneck — what is the current bottleneck, and which profile, trace or per-stage measurement shows it?
  2. Value — what does the change buy the user or the system (latency, throughput, memory, cost), and at which workload?
  3. A/B or ablation — is there a same-workload, same-head/config comparison that isolates each main claim on its own? For stacked optimizations, one number per item rather than a blended delta.

When you answer, the evidence that settles it is: base and head SHA, hardware, model, workload, warm-up and repeat count, mean or percentiles with their spread, and a correctness/quality-equivalence signal; an end-to-end claim also needs stage attribution.

@vllm-omni-review-bot

vllm-omni-review-bot commented Sep 30, 2026 •

Copy link
Copy Markdown

Omni ReviewBot: superseded

The CI failure noted on a07c91d757ca refers to an earlier head; the pull request now points at fd9e573fbd32.

Restore the per-active-request foreach history path and remove its slot-pool kernel and dedicated tests. Keep the owned NumPy PCM batching optimization and the exact-input preemption replay contract; document allocation per active request in the deployment profile.

Signed-off-by: Sy03 <1370724210@qq.com>
@Sy0307 Sy0307 added ready label to trigger buildkite CI and removed ready label to trigger buildkite CI labels Oct 1, 2026
@vllm-omni-review-bot

vllm-omni-review-bot commented Oct 1, 2026 •

Copy link
Copy Markdown

Omni ReviewBot: superseded

The CI failure noted on 6f003b86b304 refers to an earlier head; the pull request now points at 43f9bfa0186c.

@vllm-omni-review-bot vllm-omni-review-bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Omni ReviewBot review

Scan:

Category Result
Tests / verification 3 finding(s) below
Security no finding reported
Docs / comments 3 finding(s) below
Behavior / compatibility no finding reported
Correctness 1 finding(s) below

Validated:

  • [resolved] talker_first_audio vs talker_stream_decode mutual exclusion is enforced (stream_decode.py:50) and covered by test_two_stage_first_audio_remains_independent. Residual: merge lets additional_config also set talker_first_audio.
  • [resolved] connector replacement sets drop_stale_output=True (omni_scheduler_mixin.py:320) and test_connector_prompt_replacement_drops_old_frame; residual is preemption which must default False and currently has no field.
  • [resolved] at this head streaming_decoder uses per-spec hist[2,S,H,C] double-buffer not a slot-owned history pool; residual is RING=128 with sliding_window=72 (max_frames=57) covering 25-frame prime chunks.
  • [resolved] additional_config vs connector extra: connector wins on key clash; mistaken talker_stream_decode on two-stage raises. Residual: additional-only keys (e.g. code_predictor_fused) still apply on two-stage CUDA Talker.
  • [resolved] first-round drop_stale_output AttributeError: _replace_streaming_session still sets the flag at omni_ar_scheduler.py:909 when in_flight>0 so existing drain tests keep dropping — residual: test_reset:323 reads the flag on a Request that never got it
  • [claim-verified] default qwen3_tts is still two-stage (pipeline.py:30-74); qwen3_tts_fused is one-stage PCM (pipeline.py:80-104); qwen3_tts.yaml codec_left_context_frames is 72

The primary fused single-stage path is sound, but the shared OmniARScheduler change is incomplete: drop_stale_output is read on every stale frame and never initialized, so preempted in-flight requests AttributeError outside the tests that set the flag—fix by defaulting False on OmniRequest. Two fused-adjacent leaks remain: additional_config is unioned into the shared code-predictor extra helper (Omni Talker inherits fused flags), and stream_ref_context_frames treats configured 0 as 25. The new eager_pre/eager_post kernels are the production gather/scatter for fused PCM and default CUDA first-audio, yet unnamed by tests. Running prefix-cache reset is now capability-default True, so two-stage Qwen TTS Talker loses the old arch gate. CI-red questions, SO_REUSEPORT speculation, and the fused-vs-predictor rename are dropped.

Checked, no defect found:

  • tests/config/test_omni_config.py:1049 — test_sub_config_fields_match_structured_scopes includes supports_running_prefix_cache_reset in the OmniStageModelConfig field-name set, matching the declared field on OmniStageModelConfig.

Verdict: REQUEST CHANGES

CI at 6f003b86b304 (2026-10-01T21:24:12.598184+00:00): required check(s) blocking: buildkite/vllm-omni (failed). Observed Buildkite: buildkite/vllm-omni-npu-ci (failed), buildkite/vllm-omni-amd-ci (failed), buildkite/vllm-omni (failed), and 1 more.

See inline comments below.


🤖 This review was generated by InferMatrix Copilot, an open-source repo-maintenance agent for PR review, CI debugging and issue triage. Try it on your own repo, and ⭐ star it if it helped!

request.async_tokens_to_discard = max(0, stale_async_tokens - len(generated_token_ids))

if output_is_stale or async_output_is_stale:
if async_output_is_stale or (output_is_stale and request.drop_stale_output):

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P1] omni_ar_scheduler.py:537 now drops only when `async_output_is_stale or (output_…

omni_ar_scheduler.py:537 now drops only when async_output_is_stale or (output_is_stale and request.drop_stale_output). output_is_stale is True whenever num_stale_output_tokens > 0 (514–515), so a preempted in-flight frame evaluates the attribute. Unchanged by this diff: OmniRequest only inits num_stale_output_tokens (request.py:50). The only production writes are = True at connector replacement / segment stop (omni_scheduler_mixin.py:225/320, omni_ar_scheduler.py:728/909). test_preemption_stale_output_follows_upstream_delivery_policy assigns request.drop_stale_output = drop itself. A real Request that reaches this line with stale tokens and no prior True write AttributeErrors in update_from_output (OmniCFGARScheduler inherits this method). Initialize self.drop_stale_output = False next to num_stale_output_tokens and read with getattr(request, "drop_stale_output", False) — False is the preemption-delivery policy the test encodes.

Evidence: vllm_omni/core/sched/omni_ar_scheduler.py:514-515 if request is not None and request.num_stale_output_tokens > 0: then output_is_stale = True — any positive stale counter takes the new gate. vllm_omni/core/sched/omni_ar_scheduler.py:537 if async_output_is_stale or (output_is_stale and request.drop_stale_output): — unguarded attribute read when output_is_stale. unchanged by this diff, present in the PR-time tree: vllm_omni/request.py:50 self.num_stale_output_tokens = int(getattr(self, "num_stale_output_tokens", 0) or 0) — no drop_stale_output init. tests/core/sched/test_omni_ar_scheduler_stale_drain.py:257 request.drop_stale_output = drop — test injects the missing field. tests/core/sched/test_omni_ar_scheduler_logprobs.py:100-102 (unchanged by this diff) # vLLM 0.27 (a0c092ee72): Request gained num_stale_output_tokens to / # track in-flight outputs discarded at preemption/streaming-stop. / self.num_stale_output_tokens = 0 — upstream preemption uses that counter; stub has no drop_stale_output.

Suggestion: if async_output_is_stale or (output_is_stale and getattr(request, "drop_stale_output", False)):

assert scheduler.running == [request]
assert request.status == RequestStatus.RUNNING
assert request.num_in_flight_tokens == request.num_output_placeholders == 1
assert not request.drop_stale_output

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P1] Blocked reset_prefix_cache never assigns drop_stale_output, and `_make_sess…

Blocked reset_prefix_cache never assigns drop_stale_output, and _make_session/OmniRequest never initialize it (vllm_omni still patches Request to OmniRequest). assert not request.drop_stale_output therefore AttributeErrors on the blocked arm instead of proving no preemption. Same-file tests that read the flag first assign it (= drop / = False). Use getattr(request, "drop_stale_output", False) or set session.drop_stale_output = False in _make_session.

Evidence: tests/core/sched/test_omni_ar_scheduler_stale_drain.py:323 assert not request.drop_stale_output — blocked arm of this PR's new test. tests/core/sched/test_omni_ar_scheduler_stale_drain.py:49-64 _make_session builds Request(...) and only backfills num_stale_output_tokens / num_in_flight_tokens; it never sets drop_stale_output. Unchanged by this diff, present in the PR-time tree: vllm_omni/request.py:50 self.num_stale_output_tokens = int(getattr(self, "num_stale_output_tokens", 0) or 0) — OmniRequest init has no drop_stale_output default (patch.py still aliases Request to this class). vllm_omni/core/sched/omni_ar_scheduler.py:118 return False on the blocked reset path, with no drop_stale_output assign. Contrast same-file tests that assign before reading: :257 request.drop_stale_output = drop and :271 request.drop_stale_output = False.

Suggestion: assert not getattr(request, "drop_stale_output", False)

tl.store(fa_ptr + last, scheduled.to(fa_ptr.dtype.element_ty))


def eager_pre(meta, n, qsl, sampled, hidden, emb_w, ids, emb, mtp_hidden, step):

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] This PR adds eager_pre/eager_post (eager_mtp_kernels.py:104/115) as the CUD…

This PR adds eager_pre/eager_post (eager_mtp_kernels.py:104/115) as the CUDA gather/scatter around Talker-MTP. _fast_path_ok (eager_mtp.py:351) routes both fused stream-decode and the default CUDA talker_first_audio path through them when embeddings are unquantized, prefix cache is off, and text_hidden is CUDA (mtp_eager_frames = talker_first_audio_enabled(...) or self.stream_decode; qwen3_tts.yaml:27). Stream decode raises if that path is skipped (eager_mtp.py:282, not :337). tests/ never names eager_pre/eager_post. Existing tests/worker_v2/test_omni_model_state.py run_eager_mtp cases pin the CPU index_select fallback and cannot enter _fast_path_ok because text_hidden.is_cuda is false; fused-predictor and stream-audio tests never drive this indexing. Add a CUDA unit test that pins row meta, query_start_loc, sampled CB0, and scatter into codes_out/valid_out against that index_select path, including a mixed prefill/decode batch.

Evidence: eager_mtp_kernels.py:104 def eager_pre(meta, n, qsl, sampled, hidden, emb_w, ids, emb, mtp_hidden, step): — new Triton gather wrapper. eager_mtp.py:276 if self._fast_path_ok(codes_out, text_hidden): then _run_eager_mtp_fused(...). eager_mtp.py:282 if getattr(self.owner.model, "stream_decoder", None) is not None: / :283 raise RuntimeError("Talker stream decode requires the fused eager-MTP path"). eager_mtp.py:351 def _fast_path_ok(...) / :354 text_hidden.is_cuda / :355 and not self.owner.vllm_config.cache_config.enable_prefix_caching / :356 and self._embedding_weight() is not None. eager_mtp.py:424 last_tokens, layer0 = eager_pre( / :435 valid = eager_post(. qwen3_tts_talker.py:425 self.mtp_eager_frames = talker_first_audio_enabled(vllm_config) or self.stream_decode. unchanged by this diff, present in the PR-time tree: vllm_omni/deploy/qwen3_tts.yaml:27 talker_first_audio: true. tests/worker_v2/test_omni_model_state.py:468 text_hidden = torch.arange(5 * _EAGER_DIM, dtype=torch.float32).reshape(5, _EAGER_DIM) / :472 state.run_eager_mtp(...) — CPU tensors, so _fast_path_ok is false; tests/ has zero matches for eager_pre|eager_post|eager_mtp_kernels.

# options; it sets them in the stage's ``additional_config``.
additional = getattr(vllm_config, "additional_config", None)
if isinstance(additional, dict) and additional:
return {**additional, **extra_cfg}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] This hunk unions vllm_config.additional_config into shared `CodePredictorWrap…

This hunk unions vllm_config.additional_config into shared CodePredictorWrapper._stage_connector_extra_config (return {**additional, **extra_cfg} at qwen3_code_predictor.py:800). CodePredictorWrapper.__init__ then reads code_predictor_prefix_graphs / buckets / seq_lens from that helper (688–705). Unchanged by this diff: Qwen3OmniMoeTalkerCodePredictor subclasses the wrapper and calls super().__init__(vllm_config=...) (qwen3_omni_moe_code_predictor_mtp.py:12,17), so any non-empty Omni Talker additional_config that carries those keys can enable prefix re-prefill on a path that never selected them. code_predictor_fused / code_predictor_kv_cache are TTS-subclass reads of the same helper, not Omni. Keep connector extra as the shared contract; move the additional_config union into the Qwen3-TTS subclass or a TTS-only override.

Evidence: vllm_omni/model_executor/models/common/qwen3_code_predictor.py:798-800 additional = getattr(vllm_config, "additional_config", None) / if isinstance(additional, dict) and additional: / return {**additional, **extra_cfg} — shared helper now promotes any non-empty stage additional_config into predictor extra. Same file:688-694 prefix_graph_cfg = self._stage_connector_extra_config(vllm_config) / prefix_graphs_requested = self._parse_bool_config(prefix_graph_cfg.get("code_predictor_prefix_graphs")) / self._prefix_reprefill_enabled = prefix_graphs_requested and not is_npu — wrapper init arms prefix re-prefill from that merged dict. Unchanged by this diff, present in the PR-time tree: vllm_omni/model_executor/models/qwen3_omni/qwen3_omni_moe_code_predictor_mtp.py:12 class Qwen3OmniMoeTalkerCodePredictor(CodePredictorWrapper): and :17 super().__init__( with vllm_config=vllm_config — Omni inherits the new merge with no override.

Suggestion: extra_cfg = extra_cfg if isinstance(extra_cfg, dict) else {}
return extra_cfg

from .qwen3_tts_code_predictor_vllm import Qwen3TTSTalkerCodePredictorForConditionalGenerationVLLM as Predictor

extra = Predictor._stage_connector_extra_config(vllm_config)
return int(extra.get("ref_code_context_frames") or extra.get("codec_left_context_frames", 25))

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] This new resolver uses `extra.get("ref_code_context_frames") or extra.get("code…

This new resolver uses extra.get("ref_code_context_frames") or extra.get("codec_left_context_frames", 25), so a configured 0 is treated as missing and becomes 25. The fused profile and docs treat additional_config.ref_code_context_frames as the independent clone-context control (yaml sets 72). A Base/clone overlay that sets ref_code_context_frames: 0 would silently prime with 25 frames. Distinguish missing from 0 with is not None. If 0 did reach get_stream_ref_context, ref[-0:] is the full tensor; a true zero-frame bound would also need a > 0 guard like the two-stage processor.

Evidence: vllm_omni/model_executor/models/qwen3_tts/stream_decode.py:60 return int(extra.get("ref_code_context_frames") or extra.get("codec_left_context_frames", 25)) — Python or collapses 0 to the fallback 25. Consumer in this PR: vllm_omni/model_executor/models/qwen3_tts/qwen3_tts_talker.py:1516 return ref.reshape(-1, int(self.talker_config.num_code_groups))[-self.stream_ref_context_frames :] — t[-0:] is the full tensor. Unchanged two-stage sibling in the PR-time tree: vllm_omni/model_executor/stage_input_processors/qwen3_tts.py:135 ref_code_context_frames = int(cfg.get("ref_code_context_frames") or left_context_size_config) and :382 if ref_code_context_frames > 0 and int(ref_context.shape[0]) > ref_code_context_frames:.

Suggestion: extra = Predictor._stage_connector_extra_config(vllm_config)
if extra.get("ref_code_context_frames") is not None:
return int(extra["ref_code_context_frames"])
return int(extra.get("codec_left_context_frames", 25))

| `qwen3_tts_high_concurrency_mrv2_b4.yaml` | V2 | B1, B2, B3, B4 | Experimental throughput / buffered playback |
| `qwen3_tts_high_concurrency_mrv2_single_gpu.yaml` | V2 | B1–B8 | Experimental two-stage deployment on one GPU, with MPS |
| `qwen3_tts_high_concurrency.yaml` | V1 | Existing defaults | V1 high-concurrency control |
| `qwen3_tts_fused_single_gpu.yaml` | V2 | N/A (in-Talker decoder) | Opt-in single-stage CUDA pipeline on one GPU |

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] The new fused row is a V2 profile with no Code2Wav

The new fused row is a V2 profile with no Code2Wav. The following paragraph still says the V2 profiles select the Talker AR runner and the Code2Wav generation runner together. Scope that sentence to two-stage V2 profiles.

Evidence: docs/configuration/stage_configs.md:324 | qwen3_tts_fused_single_gpu.yaml | V2 | N/A (in-Talker decoder) | Opt-in single-stage CUDA pipeline on one GPU |. Unchanged by this diff, present in the PR-time tree: docs/configuration/stage_configs.md:332-334 The V2 profiles bound Talker prefill to 512 tokens per step and select the Talker AR runner and the Code2Wav generation runner together.. vllm_omni/model_executor/models/qwen3_tts/pipeline.py:77-79 # Single-stage variant: the Talker decodes each frame with the stateful streaming codec decoder and emits PCM as the final output (enable with the ``talker_stream_decode`` model option). No Code2Wav stage, no connector.

Suggestion: The two-stage V2 profiles bound Talker prefill to 512 tokens per step and select the
Talker AR runner and the Code2Wav generation
runner together. The single-stage fused profile has no Code2Wav runner.

base_config: qwen3_tts_fused_single_gpu.yaml
stages:
- stage_id: 0
additional_config:

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] The first-audio overlay in docs/configuration/stage_configs.md:386-397 restates…

The first-audio overlay in docs/configuration/stage_configs.md:386-397 restates the entire additional_config map to flip talker_stream_first_audio. resolve_deploy_yaml merges matching stages with _deep_merge_stage (stage_config.py:766), which recursively merges only _DEEP_MERGE_KEYS (stage_config.py:709-716: sampling params, engine_extras, engine_args). additional_config is not in that set, so _merge_config_fields replaces it wholesale (stage_config.py:723-728). A one-key overlay would drop talker_stream_decode; QWEN3_TTS_FUSED_PIPELINE still has engine_output_type=audio, so talker_stream_decode_enabled raises. Two-stage connectors.*.extra is different: _merge_connectors recursively merges. Add additional_config to _DEEP_MERGE_KEYS so the opt-in can be a single key, or state in this example that the map is replaced and must be copied in full.

Evidence: docs/configuration/stage_configs.md:386-397 (this diff) base_config: qwen3_tts_fused_single_gpu.yaml / additional_config: copies talker_stream_decode, talker_stream_first_audio, ref_code_context_frames, and the three code_predictor flags. Unchanged by this diff, present in the PR-time tree: stage_config.py:709-716 _DEEP_MERGE_KEYS = frozenset({"default_sampling_params", "default_pooling_params", "subtalker_sampling_params", "engine_extras", "engine_args"}) — no additional_config. stage_config.py:723-728 def _merge_config_fields(...): docstring Recursively merge selected fields; overlay replaces all other fields. and return {**base, **overlay, **merged_nested}. stage_config.py:766 by_id[sid] = _deep_merge_stage(by_id[sid], overlay_stage). Contrast, unchanged: stage_config.py:791-798 _merge_connectors return _get_recursively_merged_dict(base or {}, overlay or {}). Eligibility if the map is replaced: stream_decode.py:19-21 if not requested: / if getattr(model, "engine_output_type", None) == "audio": / raise ValueError("The single-stage Qwen3-TTS pipeline requires talker_stream_decode=True").

@vllm-omni-review-bot

Copy link
Copy Markdown

Omni ReviewBot: finding feedback

**[p1] Blocked reset_prefix_cache never assigns drop_stale_output, and _make_sess…** — tests/core/sched/test_omni_ar_scheduler_stale_drain.py:323`

Blocked reset_prefix_cache never assigns drop_stale_output, and _make_session/OmniRequest never initialize it (vllm_omni still patches Request to OmniRequest). assert not request.drop_stale_output therefore AttributeErrors on the blocked arm instead of proving no preemption. Same-file tests that read the flag first assign it (= drop / = False). Use getattr(request, "drop_stale_output", False) or set session.drop_stale_output = False in _make_session.

Evidence: tests/core/sched/test_omni_ar_scheduler_stale_drain.py:323 assert not request.drop_stale_output — blocked arm of this PR's new test. tests/core/sched/test_omni_ar_scheduler_stale_drain.py:49-64 _make_session builds Request(...) and only backfills num_stale_output_tokens / num_in_flight_tokens; it never sets drop_stale_output. Unchanged by this diff, present in the PR-time tree: vllm_omni/request.py:50 self.num_stale_output_tokens = int(getattr(self, "num_stale_output_tokens", 0) or 0) — OmniRequest init has no drop_stale_output default (patch.py still aliases Request to this class). vllm_omni/core/sched/omni_ar_scheduler.py:118 return False on the blocked reset path, with no drop_stale_output assign. Contrast same-file tests that assign before reading: :257 request.drop_stale_output = drop and :271 request.drop_stale_output = False.

Suggestion: assert not getattr(request, "drop_stale_output", False)

This finding appears in the bot's COMMENT review, but GitHub could not place it as an inline diff comment. If you are the PR author and disagree, react 👎 to this comment. The disagreement will be shown to the maintainer; it does not approve or merge the PR.

@vllm-omni-review-bot

Copy link
Copy Markdown

Omni ReviewBot: finding feedback

**[p1] omni_ar_scheduler.py:537 now drops only when async_output_is_stale or (output_…** — vllm_omni/core/sched/omni_ar_scheduler.py:537`

omni_ar_scheduler.py:537 now drops only when async_output_is_stale or (output_is_stale and request.drop_stale_output). output_is_stale is True whenever num_stale_output_tokens > 0 (514–515), so a preempted in-flight frame evaluates the attribute. Unchanged by this diff: OmniRequest only inits num_stale_output_tokens (request.py:50). The only production writes are = True at connector replacement / segment stop (omni_scheduler_mixin.py:225/320, omni_ar_scheduler.py:728/909). test_preemption_stale_output_follows_upstream_delivery_policy assigns request.drop_stale_output = drop itself. A real Request that reaches this line with stale tokens and no prior True write AttributeErrors in update_from_output (OmniCFGARScheduler inherits this method). Initialize self.drop_stale_output = False next to num_stale_output_tokens and read with getattr(request, "drop_stale_output", False) — False is the preemption-delivery policy the test encodes.

Evidence: vllm_omni/core/sched/omni_ar_scheduler.py:514-515 if request is not None and request.num_stale_output_tokens > 0: then output_is_stale = True — any positive stale counter takes the new gate. vllm_omni/core/sched/omni_ar_scheduler.py:537 if async_output_is_stale or (output_is_stale and request.drop_stale_output): — unguarded attribute read when output_is_stale. unchanged by this diff, present in the PR-time tree: vllm_omni/request.py:50 self.num_stale_output_tokens = int(getattr(self, "num_stale_output_tokens", 0) or 0) — no drop_stale_output init. tests/core/sched/test_omni_ar_scheduler_stale_drain.py:257 request.drop_stale_output = drop — test injects the missing field. tests/core/sched/test_omni_ar_scheduler_logprobs.py:100-102 (unchanged by this diff) # vLLM 0.27 (a0c092ee72): Request gained num_stale_output_tokens to / # track in-flight outputs discarded at preemption/streaming-stop. / self.num_stale_output_tokens = 0 — upstream preemption uses that counter; stub has no drop_stale_output.

Suggestion: if async_output_is_stale or (output_is_stale and getattr(request, "drop_stale_output", False)):

This finding appears in the bot's COMMENT review, but GitHub could not place it as an inline diff comment. If you are the PR author and disagree, react 👎 to this comment. The disagreement will be shown to the maintainer; it does not approve or merge the PR.

Sy0307 added 7 commits October 2, 2026 05:38
Give col2im several rows per program and pick im2col, SnakeBeta and conv_out tile sizes from a bit-exact H200 sweep; conv_out drops from 256 to 32 rows per program, which had needed 255 registers per thread. Only launch shapes change and every output element keeps its operations and order, so PCM is unchanged while one batch-96 decode call falls from 2463 to 2141 us.

Signed-off-by: Sy03 <1370724210@qq.com>
Upload the text ids of a step's new non-streaming CustomVoice and VoiceDesign requests in one pinned copy and run the text embedding and projection once over their concatenated tokens; build_prompt_embeds consumes each request's rows by request id. Cache the CustomVoice speaker id and codec embedding per speaker instead of rebuilding them for every request.

Signed-off-by: Sy03 <1370724210@qq.com>
Find settled decode rows with one NumPy pass in run_preprocess instead of a per-row Python check, falling back to the row loop while a preemption replay is pending. Build the eager MTP row metadata from one transpose of the entry tuples and update ready slots in one call.

Signed-off-by: Sy03 <1370724210@qq.com>
Exercise final CustomVoice and VoiceDesign prompt assembly with batched text projection, including cache consumption and reused request ids. Compare vectorized mixed-row preprocessing with the scalar path and verify pending replay bypasses the split. Check tiled col2im output and full ping-pong state against an independent CPU oracle for reordered slots, partial row tiles and channel tails.

Signed-off-by: Sy03 <1370724210@qq.com>
Restore the original per-row settled preprocessing and eager MTP metadata
packing. Native C128 checks showed worse first-audio latency for the third
host optimization, and reverting only its NumPy split did not reliably
recover the baseline latency. Remove its unused replay query and tests.
Keep exact-input preemption replay, decoder retiling, batched prompt text
projection, and their strengthened regression tests.

Signed-off-by: Sy03 <1370724210@qq.com>
Matched native EN1088 checks do not reproduce the claimed C128 throughput
benefit, while first-audio latency increases across independent starts.
Keep the original prompt preparation and remove the batching-only tests.
Retain decoder retiling and its independent state and boundary oracle.

Signed-off-by: Sy03 <1370724210@qq.com>
Repeated native EN1088 C128 runs do not turn the isolated decoder retiling
benefit into an end-to-end throughput gain on the validated H200 setup.
Restore the original launch shapes and adapt the independent nonzero
col2im oracle to them. Production code now matches the original PR head;
the remaining increment protects PCM values, reordered slots, channel
tails, and the complete ping-pong state against an independent CPU result.

Signed-off-by: Sy03 <1370724210@qq.com>
@Sy0307 Sy0307 added ready label to trigger buildkite CI and removed ready label to trigger buildkite CI labels Oct 2, 2026

@vllm-omni-review-bot vllm-omni-review-bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Omni ReviewBot review

Scan:

Category Result
Tests / verification 5 finding(s) below
Security no finding reported
Docs / comments 2 finding(s) below
Behavior / compatibility 2 finding(s) below
Correctness 1 finding(s) below

Validated:

  • [resolved] talker_first_audio vs talker_stream_decode mutex at stream_decode.py:50; residual is the shared additional_config merge feeding both readers.
  • [resolved] first_audio vs stream_decode exclusion: stream_decode.py:50 raises mutually exclusive; residual is additional_config merge still delivering talker_first_audio, which the same function rejects.
  • [resolved] connector replacement sets drop_stale_output=True (omni_scheduler_mixin.py:320); residual blocked-reset AttributeError already in ReviewBot thread.
  • [resolved] running reset capability: topology False is forced in _build_engine_args/omni_config and test_single_stage_running_reset_capability_cannot_be_overridden; residual is getattr(..., True) in reset_prefix_cache.
  • [resolved] bot p1 drop_stale_output still_affected: OmniRequest still does not init the flag (request.py:50); line 537 still reads it without getattr — duplicate, not re-published.
  • [resolved] architecture-named running-reset block replaced by the flag (omni_ar_scheduler.py:109); residual: two-stage Talker still defaults True so async_omni.reset_prefix_cache(reset_running_requests=True) now calls super() on stock qwen3_tts — first_frame_decoder is single-frame/no slot, not published

Keep the primary PCM length-end defect (decode computed+scheduled+1 fires one token early, AudioRequest.active=False then drops later frames; +1 is only valid on prefill) and the untested fused eager-MTP kernels that are now the default CUDA path for two-stage first-audio as well as fused PCM. Also keep fused poller-default omission, CodePredictorWrapper additional_config blast radius onto Omni MoE, unseeded MTP generator cache lifecycle, ref_code_context 0→25, the untested mm_outputs_fresh_per_step snapshot skip, and a nit to fix examples that still call qwen3_tts.yaml the only deploy config. Drop already-raised drop_stale_output comments, the intentional SO_REUSEPORT scope ask, and the PCM encode no-issue note.

Checked, no defect found:

  • tests/config/test_omni_config.py:1049 — test_sub_config_fields_match_structured_scopes stays in lockstep with OmniStageModelConfig.supports_running_prefix_cache_reset (default True). QWEN3_TTS_FUSED_PIPELINE still sets topology False, and _build_model_config clamps kwargs to False when topology is False.

Verdict: REQUEST CHANGES

CI at 43f9bfa0186c (2026-10-02T08:19:11.729583+00:00): required check(s) blocking: buildkite/vllm-omni (pending). Observed Buildkite: buildkite/vllm-omni-intel-ci (passed), buildkite/vllm-omni-npu-ci (failed), buildkite/vllm-omni-amd-ci (failed), and 1 more.

Findings

  • **[P1] CUDA eager-MTP with prefix cache off and unquantized embeddings now takes _fas…** — vllm_omni/worker_v2/model_states/eager_mtp.pyCUDA eager-MTP with prefix cache off and unquantized embeddings now takes_fast_path_ok→_run_eager_mtp_fused/eager_pre/eager_postand returns, so the Python gather/scatter is skipped. Stockqwen3_tts.yaml is in that set (talker_first_audio: true, enable_prefix_caching: false); mtp_eager_framesis also set by two-stage first-audio, and fused decode writesfirst_audioviaHAS_FA_VALIDwithoutstream_decoder. Existing first-audio tests in test_omni_model_state.pyuse CPUtorch.zeros, so they never enter this path. No test names eager_pre, eager_post, _fast_path_ok, or _run_eager_mtp_fused. Add a two-stage first-audio case that fails if the fused path is not taken and pins last_tokens/layer0/HAS_FA_VALIDagainst the Python path. Do not gate onstream_decoder`.

Evidence: eager_mtp.py:276 if self._fast_path_ok(codes_out, text_hidden): then _run_eager_mtp_fused(...) and return — Python gather/scatter is not used. eager_mtp.py:351-359 _fast_ok = (text_hidden.is_cuda and not self.owner.vllm_config.cache_config.enable_prefix_caching and self._embedding_weight() is not None and hasattr(self.owner.model, "_codebook_vocab_size")) with no stream_decoder gate. eager_mtp.py:282-283 if getattr(self.owner.model, "stream_decoder", None) is not None: raise RuntimeError("Talker stream decode requires the fused eager-MTP path") — fused PCM requires this path; two-stage first-audio still takes it when _fast_path_ok is true. eager_mtp.py:433-438 fa_in_kernel = isinstance(first_audio, torch.Tensor) and not has_prefill then eager_post(..., first_audio if fa_in_kernel else None, self._first_audio_valid if fa_in_kernel else None, ...) — two-stage first-audio decode writes the marker in-kernel. eager_mtp_kernels.py:95-101 if WRITE_FA: scheduled = tl.load(meta_ptr + 3 * n + b) != 0 / if HAS_FA_VALID: scheduled = scheduled & (tl.load(fa_valid_ptr + req) != 0) / else: scheduled = scheduled & False. qwen3_tts_talker.py:425 self.mtp_eager_frames = talker_first_audio_enabled(vllm_config) or self.stream_decode. Unchanged by this diff, present in the PR-time tree: qwen3_tts.yaml:27 talker_first_audio: true; qwen3_tts.yaml:80 enable_prefix_caching: false. Unchanged by this diff, present in the PR-time tree: tests/worker_v2/test_omni_model_state.py:574 state.run_eager_mtp(batch, torch.zeros(1, _EAGER_DIM), torch.tensor([[_EOS]]), outputs) and :597 state.run_eager_mtp(batch, torch.zeros(2, _EAGER_DIM), torch.tensor([[7], [8]]), outputs) — CPU torch.zeros, so _fast_path_ok is false. Repo-wide tests/ grep for eager_pre|eager_post|_fast_path_ok|_run_eager_mtp_fused returns no matches.

Suggestion: if self._fast_path_ok(codes_out, text_hidden) and getattr(self.owner.model, "stream_decoder", None) is not None:

  • **[P2] stream_ref_context_frames uses extra.get("ref_code_context_frames") or extra.g…** — vllm_omni/model_executor/models/qwen3_tts/stream_decode.pystream_ref_context_frames usesextra.get("ref_code_context_frames") or extra.get("codec_left_context_frames", 25), so a YAML overlay of 0 is falsy and becomes 25. get_stream_ref_context always slices [-n:]; Python [-0:]is the full tensor, and empty/missingcodes.ref` already returns None, so 0 would not mean “no clone context” even after a None-check. Use an explicit None check so 0 is a real boundary; if 0 should skip priming, return None when frames <= 0 instead of slicing.

Evidence: stream_decode.py:60 return int(extra.get("ref_code_context_frames") or extra.get("codec_left_context_frames", 25)) — configured 0 is falsy and becomes 25. qwen3_tts_talker.py:1514-1516 if not isinstance(ref, torch.Tensor) or not ref.numel(): / return None / return ref.reshape(-1, int(self.talker_config.num_code_groups))[-self.stream_ref_context_frames :] — missing/empty ref already skips clone priming; [-0:] would keep the full tensor, not an empty slice.

Suggestion: frames = extra.get("ref_code_context_frames")
if frames is None:
frames = extra.get("codec_left_context_frames", 25)
return int(frames)

  • [P2] Talker mm_outputs_fresh_per_step (qwen3_tts_talker.py:738-740) is the only… — vllm_omni/model_executor/models/qwen3_tts/qwen3_tts_talker.py
    Talker mm_outputs_fresh_per_step (qwen3_tts_talker.py:738-740) is the only gate that skips _retain_multimodal_outputs snapshotting (omni_ar_model_runner.py:310 getattr). No test names that property. test_async_mm_snapshot_owns_output_until_copy_finishes uses runner.model = SimpleNamespace(), so getattr defaults False and never takes the skip; test_stream_audio_snapshot_skips_discarded_code_partition stubs _build_async_chunk_outputs_from_mm (a later partition helper) and never sets the property. Renaming the attribute would silently restore snapshotting. Pin skip when the property is True (stream_decoder set and eager_frames_active) and snapshot when it is False.

Evidence: qwen3_tts_talker.py:738-740 @property / def mm_outputs_fresh_per_step(self) -> bool: / return self.stream_decoder is not None and self.eager_frames_active. omni_ar_model_runner.py:310 if getattr(self.model, "mm_outputs_fresh_per_step", False): — skip snapshot. tests/worker_v2/test_omni_ar_model_runner.py:144 runner.model = SimpleNamespace() — existing _retain_multimodal_outputs test has no property, getattr defaults False. tests/worker_v2/test_omni_ar_model_runner.py:650-654 monkeypatch.setattr(OmniARModelRunner, "_build_async_chunk_outputs_from_mm", lambda *args: pytest.fail(...)) — streaming test stubs a different helper and never sets mm_outputs_fresh_per_step. Repo-wide grep of mm_outputs_fresh_per_step hits only those two production sites (no test file).


🤖 This review was generated by InferMatrix Copilot, an open-source repo-maintenance agent for PR review, CI debugging and issue triage. Try it on your own repo, and ⭐ star it if it helped!

try:
if not sock.getsockopt(socket.SOL_SOCKET, socket.SO_REUSEPORT):
return sock
own = socket.socket(family=sock.family, type=socket.SOCK_STREAM)

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] _own_reuseport_socket creates own at line 269, then except OSError return…

_own_reuseport_socket creates own at line 269, then except OSError returns the inherited sock without own.close(), and the success path return own never closes sock. run_omni_api_server_worker_proc (line 308) calls this for every Omni HTTP frontend when client_count>1. Close own on the failure path and close the inherited socket after a successful bind.

Evidence: vllm_omni/entrypoints/openai/api_server.py:269 own = socket.socket(family=sock.family, type=socket.SOCK_STREAM) — created inside the try; api_server.py:273-276 own.bind(sock.getsockname()[:2]) / except OSError as e: / return sock / return own — bind failure returns inherited sock with no own.close(), success returns own with no sock.close(); api_server.py:308 sock = _own_reuseport_socket(sock, int(omni_client_config["client_count"])) — every Omni API worker with client_count>1. Unchanged by this diff, present in the PR-time tree: vllm_omni/entrypoints/cli/serve.py:1265 listen_address, sock = setup_server(args, reuse_port=True) and serve.py:1285 target_server_fn=run_omni_api_server_worker_proc — multi-frontend path that feeds that worker.

self._fused = FusedCodePredictor(self, max_batch)
# Compile every kernel variant outside any graph capture.
hidden = int(self.talker_config.hidden_size)
for batch in range(1, min(2, max_batch) + 1):

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] Fused _setup_compile comments "Compile every kernel variant" then only runs `…

Fused _setup_compile comments "Compile every kernel variant" then only runs _fused(..., top_k=50) for batch in range(1, min(2, max_batch)+1). Talker configure_mtp_execution_buckets is meant to cover every reachable MRv2 size for FrameLocalKV (for batch in self._bucket_sizes); fused warmup ignores those buckets. _cp_sample_kernel takes TOPK as tl.constexpr from live int(top_k), so a non-50 request JITs on first use. Batch is only the launch grid, so sweeping buckets would not compile extra kernels. Warm the live top_k, not only 50.

Evidence: qwen3_tts_code_predictor_vllm.py:106-114 # Compile every kernel variant outside any graph capture. / for batch in range(1, min(2, max_batch) + 1): / self._fused(codes, zeros, zeros, 1.0, 50, uniforms) — fused warmup hardcodes top_k=50 and only batch 1–2. qwen3_tts_code_predictor_vllm.py:148-154 return self._fused(..., int(top_k), sample_uniforms) — live top_k is forwarded. fused_code_predictor.py:157 TOPK: tl.constexpr, and fused_code_predictor.py:304-307 _cp_sample_kernel[(B,)]( / V=self.vocab, TOPK=top_k, ... — TOPK specializes the sample kernel; B is only the grid. Present in the PR-time tree (this region is not the fused insert): qwen3_tts_talker.py:539-548 # Align predictor warmup/capture buckets with the outer MRv2 MTP / # graph buckets so the compile cache covers every reachable size. / self.code_predictor.configure_mtp_execution_buckets(sorted(outer_buckets)). qwen3_tts_fused_single_gpu.yaml:46 and :53 top_k: 50 — default and subtalker sampling both use 50.

# streaming codec decoder and emits PCM as the final output (enable with the
# ``talker_stream_decode`` model option). No Code2Wav stage, no connector.
QWEN3_TTS_FUSED_PIPELINE = PipelineConfig(
model_type="qwen3_tts_fused",

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] This PR registers pipeline model_type="qwen3_tts_fused" (pipeline.py:81)

This PR registers pipeline model_type="qwen3_tts_fused" (pipeline.py:81). _event_driven_orch_default_for_pipeline still exact-matches "qwen3_tts" (orchestrator.py:104, unchanged), so fused serving keeps the legacy 1 ms poll unless VLLM_OMNI_EVENT_DRIVEN_ORCH is set. test_only_qwen3_tts_has_a_pipeline_default (test_orchestrator_event_driven.py:131-134) lists other types as False and omits fused. Include qwen3_tts_fused in that default and the test.

Evidence: pipeline.py:81 model_type="qwen3_tts_fused", registers the new pipeline. Unchanged by this diff, present in the PR-time tree: orchestrator.py:87 # Default is off except for pipelines with an explicit validated default.; orchestrator.py:104 return pipeline_model_type == "qwen3_tts" — exact-match only, so fused is False. Unchanged wiring: omni_engine_base.py:886-888 self._event_driven_orch_default = _event_driven_orch_default_for_pipeline(pipeline_config.model_type if pipeline_config is not None else None). Unchanged by this diff: tests/engine/test_orchestrator_event_driven.py:132-134 assert _event_driven_orch_default_for_pipeline("qwen3_tts") is True then for model_type in (None, "qwen3_omni_moe", "minicpmo_4_5", "moss_tts_delay"): with no fused.

# options; it sets them in the stage's ``additional_config``.
additional = getattr(vllm_config, "additional_config", None)
if isinstance(additional, dict) and additional:
return {**additional, **extra_cfg}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

**[P2] CodePredictorWrapper._stage_connector_extra_config now returns `{additional,…

CodePredictorWrapper._stage_connector_extra_config now returns {**additional, **extra_cfg} (qwen3_code_predictor.py:800) so single-stage TTS can read additional_config without a connector. That helper is shared: wrapper init already consumes it for code_predictor_prefix_graphs / _buckets / _seq_lens (same file:688-705), and unchanged Qwen3OmniMoeTalkerCodePredictor subclasses the wrapper and calls super().init. Omni talker additional_config/CLI overlays therefore alias those connector extras. Keep connector extras in extra and merge additional_config only in the Qwen3-TTS extra helper.

Evidence: qwen3_code_predictor.py:798-800 additional = getattr(vllm_config, "additional_config", None) / if isinstance(additional, dict) and additional: / return {**additional, **extra_cfg} — shared helper now folds stage additional_config into connector extras. Existing consumer in the same class (not introduced by this hunk, present in the PR-time tree): qwen3_code_predictor.py:688-705 prefix_graph_cfg = self._stage_connector_extra_config(vllm_config) then prefix_graph_cfg.get("code_predictor_prefix_graphs"), code_predictor_prefix_graph_buckets, code_predictor_prefix_graph_seq_lens. Unchanged by this diff, present in the PR-time tree: qwen3_omni_moe_code_predictor_mtp.py:12 class Qwen3OmniMoeTalkerCodePredictor(CodePredictorWrapper): and :17-18 super().__init__(vllm_config=vllm_config, ...).

Suggestion: extra_cfg = extra_cfg if isinstance(extra_cfg, dict) else {}
return extra_cfg

cache = {}
self._mtp_generators = cache
generator = cache.get(req_id, _NO_SEED)
if generator is None:

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] _get_mtp_generator stores cache[req_id] = None when get_mtp_seed returns…

_get_mtp_generator stores cache[req_id] = None when get_mtp_seed returns None (omni_model_state.py:1202-1204) and later cache.get(req_id, _NO_SEED) treats that None as final (1194-1197). New finish_audio pops _stream_pos/_talker_inputs but not _mtp_generators (eager_mtp.py:105-110); remove_request only runs when the runner slot is still mapped. Same-id reuse after finish then skips get_mtp_seed — the seed test already has to clear() before treating r1 as a new request. Do not cache None, and pop generators in finish_audio.

Evidence: omni_model_state.py:1194-1197 generator = cache.get(req_id, _NO_SEED) / if generator is None: # Resolved earlier: this request has no seed. return None; omni_model_state.py:1202-1204 if seed is None: cache[req_id] = None; return None. eager_mtp.py:105-110 (new in this PR) def finish_audio pops _stream_pos and _talker_inputs only — no _mtp_generators.pop. Unchanged by this diff, present in the PR-time tree: omni_model_runner.py:764-767 for req_id in all_done: idx = self.req_states.req_id_to_index.get(req_id); if idx is not None: self.model_state.remove_request(idx) — generator pop is skipped when the slot is already released. tests/worker_v2/test_omni_model_state.py:272 state._mtp_generators.clear() # The following rows are new requests.

Suggestion: if generator is None:
return None

@vllm-omni-review-bot

Copy link
Copy Markdown

Omni ReviewBot: finding feedback

**[p1] CUDA eager-MTP with prefix cache off and unquantized embeddings now takes _fas…** — vllm_omni/worker_v2/model_states/eager_mtp.py`

CUDA eager-MTP with prefix cache off and unquantized embeddings now takes _fast_path_ok → _run_eager_mtp_fused/eager_pre/eager_post and returns, so the Python gather/scatter is skipped. Stock qwen3_tts.yaml is in that set (talker_first_audio: true, enable_prefix_caching: false); mtp_eager_frames is also set by two-stage first-audio, and fused decode writes first_audio via HAS_FA_VALID without stream_decoder. Existing first-audio tests in test_omni_model_state.py use CPU torch.zeros, so they never enter this path. No test names eager_pre, eager_post, _fast_path_ok, or _run_eager_mtp_fused. Add a two-stage first-audio case that fails if the fused path is not taken and pins last_tokens/layer0/HAS_FA_VALID against the Python path. Do not gate on stream_decoder.

Evidence: eager_mtp.py:276 if self._fast_path_ok(codes_out, text_hidden): then _run_eager_mtp_fused(...) and return — Python gather/scatter is not used. eager_mtp.py:351-359 _fast_ok = (text_hidden.is_cuda and not self.owner.vllm_config.cache_config.enable_prefix_caching and self._embedding_weight() is not None and hasattr(self.owner.model, "_codebook_vocab_size")) with no stream_decoder gate. eager_mtp.py:282-283 if getattr(self.owner.model, "stream_decoder", None) is not None: raise RuntimeError("Talker stream decode requires the fused eager-MTP path") — fused PCM requires this path; two-stage first-audio still takes it when _fast_path_ok is true. eager_mtp.py:433-438 fa_in_kernel = isinstance(first_audio, torch.Tensor) and not has_prefill then eager_post(..., first_audio if fa_in_kernel else None, self._first_audio_valid if fa_in_kernel else None, ...) — two-stage first-audio decode writes the marker in-kernel. eager_mtp_kernels.py:95-101 if WRITE_FA: scheduled = tl.load(meta_ptr + 3 * n + b) != 0 / if HAS_FA_VALID: scheduled = scheduled & (tl.load(fa_valid_ptr + req) != 0) / else: scheduled = scheduled & False. qwen3_tts_talker.py:425 self.mtp_eager_frames = talker_first_audio_enabled(vllm_config) or self.stream_decode. Unchanged by this diff, present in the PR-time tree: qwen3_tts.yaml:27 talker_first_audio: true; qwen3_tts.yaml:80 enable_prefix_caching: false. Unchanged by this diff, present in the PR-time tree: tests/worker_v2/test_omni_model_state.py:574 state.run_eager_mtp(batch, torch.zeros(1, _EAGER_DIM), torch.tensor([[_EOS]]), outputs) and :597 state.run_eager_mtp(batch, torch.zeros(2, _EAGER_DIM), torch.tensor([[7], [8]]), outputs) — CPU torch.zeros, so _fast_path_ok is false. Repo-wide tests/ grep for eager_pre|eager_post|_fast_path_ok|_run_eager_mtp_fused returns no matches.

Suggestion: if self._fast_path_ok(codes_out, text_hidden) and getattr(self.owner.model, "stream_decoder", None) is not None:

This finding appears in the bot's COMMENT review, but GitHub could not place it as an inline diff comment. If you are the PR author and disagree, react 👎 to this comment. The disagreement will be shown to the maintainer; it does not approve or merge the PR.

@linyueqian linyueqian 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.

Thanks for addressing everything. I re-ran this head (43f9bfa) on one H200-class GPU, fresh server per condition, 3 runs per cell.

  • The replacement fence now sets drop_stale_output, and the new test fails if that line is removed, so it does cover the case. The shared scheduler checks the supports_running_prefix_cache_reset capability with no model name, and reference priming moved into model hooks.
  • talker_stream_decode on the two-stage YAML now fails at model init with a clear message instead of serving empty audio.
  • Fused CustomVoice streaming: c1 15.2 audio-s/s at 14 ms TTFP, c8 104.4 at 19 ms, c32 286.2 at 29 ms, against 110.5 at 63 ms for two-stage main at c32. Non-streaming and the earlier head match within noise.
  • Quality: CustomVoice WER 2.53% vs 2.95% for two-stage main (CI [-1.11, +0.22] pp). Base with the 72-frame reference context is 2.56% vs 2.51% on main (CI [-0.28, +0.35] pp), and TTFP cost of the longer context is small (88 ms vs 77 ms at c32).
  • Length termination still flushes exactly N x 1920 samples, chunk boundaries are no worse than other frame edges, and a seeded request is identical before and after cancels and slot reuse.
  • Two-stage path on this head: c32 113.2 audio-s/s vs 110.5 on main, WER within noise.
  • The 14 touched test files pass on GPU apart from the unrelated test_diffusion_stage_payload_keys_roundtrip cases, which need HF downloads and fail the same way on main offline.

One correction to my earlier comment, for the record: I said the two-stage path was bit-identical to main at batch 1. That held for the single request I compared, but main's two-stage output is not stable across server restarts (only 9 of 55 seeded requests match between two main servers), so bit identity is not a meaningful bar there. When both servers land in the same startup state this head matches main 55 of 55. That nondeterminism is on main and not from this PR.

Looks good to me.

@Sy0307
Sy0307 merged commit e84014d into vllm-project:main Oct 2, 2026
7 of 9 checks passed
BeatSeat added a commit to BeatSeat/vllm-omni that referenced this pull request Oct 2, 2026
…elta

The API MainThread and the orchestrator thread share one GIL. At 80
sessions they handle about 1000 input frames and 200 audio deltas per
second.

- The resume journal keeps the JSON text it encodes to count bytes, and
  the websocket sends that text; the payload is no longer encoded twice.
- An input_audio_buffer.append is base64-decoded once, and the base64 text
  is kept for the session payload instead of being re-encoded.
- A Realtime audio delta updates only the marks it adds, not a re-sort of
  all marks and two rebuilds of the item's audio part (O(n) per delta).
- PCM and WAV output audio are written with numpy, byte for byte as
  soundfile writes them. A start-up probe compares the bytes and falls
  back to soundfile on any mismatch. The probe includes libsndfile 1.2's
  quantizer (x * 2**31 rounded half to even into an int32, then its high
  16 bits); without it the fast path never engaged on libsndfile 1.2.
  It also includes AudioMixin's own RAW writer (floor(x * 32768) in
  float64), which main uses for pcm since vllm-project#8259; the probe compares with
  that fallback, so the fast path reproduces whichever one is in use.

A800, N=80: 0/80 -> 80/80 sessions (RTF 1.215 -> 1.045) from the PCM_16
quantizer alone; N=64 64/64 either way.

Signed-off-by: BeatSeat <wendavid552@gmail.com>
@NickCao

NickCao commented Oct 2, 2026

Copy link
Copy Markdown
Collaborator

This breaks test-ready, see #8426

@Sy0307

Sy0307 commented Oct 2, 2026

Copy link
Copy Markdown
Collaborator Author

This breaks test-ready, see #8426

Sorry for the mistake. Have approved it and we should merge it asap.

BeatSeat added a commit to BeatSeat/vllm-omni that referenced this pull request Oct 4, 2026
…elta

The API MainThread and the orchestrator thread share one GIL. At 80
sessions they handle about 1000 input frames and 200 audio deltas per
second.

- The resume journal keeps the JSON text it encodes to count bytes, and
  the websocket sends that text; the payload is no longer encoded twice.
- An input_audio_buffer.append is base64-decoded once, and the base64 text
  is kept for the session payload instead of being re-encoded.
- A Realtime audio delta updates only the marks it adds, not a re-sort of
  all marks and two rebuilds of the item's audio part (O(n) per delta).
- PCM and WAV output audio are written with numpy, byte for byte as
  soundfile writes them. A start-up probe compares the bytes and falls
  back to soundfile on any mismatch. The probe includes libsndfile 1.2's
  quantizer (x * 2**31 rounded half to even into an int32, then its high
  16 bits); without it the fast path never engaged on libsndfile 1.2.
  It also includes AudioMixin's own RAW writer (floor(x * 32768) in
  float64), which main uses for pcm since vllm-project#8259; the probe compares with
  that fallback, so the fast path reproduces whichever one is in use.

A800, N=80: 0/80 -> 80/80 sessions (RTF 1.215 -> 1.045) from the PCM_16
quantizer alone; N=64 64/64 either way.

Signed-off-by: BeatSeat <wendavid552@gmail.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

high priority high priority issue, needs to be done asap merge-test label to trigger buildkite merge test CI ready label to trigger buildkite CI tts code related to tts models

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants