Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
3840cd1
refactor(diffusion): core data type changes for single-request contract
yJader Jun 2, 2026
b274a94
refactor(diffusion): infrastructure for request-level batching
yJader Jun 2, 2026
ad0db31
refactor(diffusion): migrate pipelines to RequestBatch forward
yJader Jun 2, 2026
5e803a4
refactor(diffusion): add request-mode batch execution
yJader Jun 2, 2026
f56f853
feat(diffusion/ipc): enhance SHM packing to support batch runner outputs
yJader Jun 3, 2026
2b13196
refactor(diffusion): determine request batching support during engine…
yJader Jun 4, 2026
a595d61
style: simplify code annotation
yJader Jun 4, 2026
2fcb296
refactor(diffusion/pipelines): migrate to RequestBatch for uncovered …
yJader Jun 7, 2026
632aa4b
style(diffusion): rename RequestBatch to DiffusionRequestBatch
yJader Jun 7, 2026
8af7a9e
docs(diffusion): update relevant documents
yJader Jun 9, 2026
ed789cb
test(diffusion): align batching tests with single-prompt requests
yJader Jun 9, 2026
9c2348e
feat(diffusion): add request batch admission window
SamitHuang Jun 15, 2026
f3486fb
feat(memory): add peak memory tracking for stepwise requests
yJader Jun 15, 2026
f349fac
docs(diffusion): add request-level batching documentation and update …
yJader Jun 15, 2026
4608e39
fix(diffusion): resolve rebase fallout in request batching and step h…
yJader Jun 16, 2026
7f1dac7
refactor(diffusion): unify prompt handling in request and output form…
yJader Jun 16, 2026
007e32a
fix: update engine configuration and fix prompt variable naming
yJader Jun 16, 2026
fae81e3
refactor: unify request handling and update prompt structure across d…
yJader Jun 17, 2026
05b0977
fix: GLM image
yJader Jun 21, 2026
c17a375
refactor: Refactor pipeline return types from list[DiffusionOutput] t…
yJader Jun 25, 2026
476c336
refactor(diffusion): roll back request-batch pipeline interface
yJader Jun 27, 2026
8d1db92
refactor(diffusion): roll back request-batch hunyuan-image pipeline i…
yJader Jun 28, 2026
c75c421
Merge upstream/main into request-batch
yJader Jun 28, 2026
03fcbdb
feat(diffusion): Enhance request batching and sampling parameter hand…
yJader Jun 28, 2026
1203d15
refactor(test): Replace SimpleNamespace with DiffusionOutput in batch…
yJader Jun 28, 2026
176544c
refactor(diffusion): Remove unsupported list-prompt batch request han…
yJader Jun 28, 2026
aa27bb3
refactor(diffusion): Simplify diffusion request-batch capability checks
yJader Jun 28, 2026
437d923
refactor(diffusion): Remove legacy diffusion batch message handler
yJader Jun 28, 2026
089a9a8
refactor(diffusion): Share diffusion request-batch output splitting
yJader Jun 28, 2026
534475b
refactor(diffusion): Unify diffusion model runner request execution
yJader Jun 28, 2026
14f14b4
fix(tests): Update OmniDiffusionRequest to use single prompt parameter
yJader Jun 29, 2026
a9c1b96
fix(diffusion/flux): Enable request batching support and update retur…
yJader Jun 29, 2026
9ffb55c
refactor(pipelines): Update request handling to support single prompt…
yJader Jun 29, 2026
dcf55c0
fix(diffusion): apply batch peak memory to every request
yJader Jun 29, 2026
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
2 changes: 2 additions & 0 deletions docs/.nav.yml
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,7 @@ nav:
- CPU Offloading: user_guide/diffusion/cpu_offload_diffusion.md
- LoRA: user_guide/diffusion/lora.md
- Custom Pipeline: features/custom_pipeline.md
- Request-Level Batching: user_guide/diffusion/request_batching.md
- Step Execution: user_guide/diffusion/step_execution.md
- Quantization:
- Overview: user_guide/quantization/overview.md
Expand Down Expand Up @@ -118,6 +119,7 @@ nav:
- design/feature/async_chunk.md
- design/feature/vae_parallel.md
- design/feature/diffusion_step_execution.md
- design/feature/diffusion_request_level_batching.md
- design/feature/diffusion_continuous_batching.md
- Module Design:
- design/module/ar_module.md
Expand Down
14 changes: 14 additions & 0 deletions docs/configuration/stage_configs.md
Original file line number Diff line number Diff line change
Expand Up @@ -379,6 +379,20 @@ The maximum number of sequences for concurrent processing in this stage. For LLM

Default: `1`

#### `engine_args.request_batch_max_wait_ms`

The maximum time, in milliseconds, that a diffusion request-mode stage may wait
before the first `schedule()` of a new scheduler wave so compatible requests can
accumulate for request-level batching. This only applies to diffusion pipelines
that support request-level batching with `step_execution` disabled.

Use this together with `max_num_seqs > 1` for bursty serving traffic. `0`
disables admission waiting and preserves the lowest first-request latency.
For diffusion request-level batching tuning, see
[Request-Level Batching](../user_guide/diffusion/request_batching.md).

Default: `0.0`

### `engine_args`

Engine arguments for configuring the LLM engine, diffusion engine, or other engine types used by this stage.
Expand Down
51 changes: 27 additions & 24 deletions docs/contributing/model/adding_diffusion_model.md
Original file line number Diff line number Diff line change
Expand Up @@ -349,12 +349,14 @@ class YourModelPipeline(nn.Module):
- def __call__(
+ def forward(
self,
+ req: OmniDiffusionRequest, # ← Add request parameter here
+ req: DiffusionRequestBatch, # ← Add request-batch parameter here
- ):
+ ) -> DiffusionOutput: # ← Add return type
+ ) -> list[DiffusionOutput]: # ← Add return type
```

[`OmniDiffusionRequest`](https://docs.vllm.ai/projects/vllm-omni/en/latest/api/vllm_omni/diffusion/request/#vllm_omni.diffusion.request.OmniDiffusionRequest) is a dataclass that contains the **prompts** and **sampling parameters** [`OmniDiffusionSamplingParams`](https://docs.vllm.ai/projects/vllm-omni/en/latest/api/vllm_omni/inputs/data/#vllm_omni.inputs.data.OmniDiffusionSamplingParams) for the diffusion pipeline execution. It also contains a request_id for other components to trace this request and its outputs.
[`OmniDiffusionRequest`](https://docs.vllm.ai/projects/vllm-omni/en/latest/api/vllm_omni/diffusion/request/#vllm_omni.diffusion.request.OmniDiffusionRequest) is a dataclass that contains one **prompt** and the **sampling parameters** [`OmniDiffusionSamplingParams`](https://docs.vllm.ai/projects/vllm-omni/en/latest/api/vllm_omni/inputs/data/#vllm_omni.inputs.data.OmniDiffusionSamplingParams) for one logical diffusion request. It also contains a request_id for other components to trace this request and its outputs. Before pipeline execution, the runner wraps one or more independent requests into `DiffusionRequestBatch`.

[`DiffusionRequestBatch`](https://docs.vllm.ai/projects/vllm-omni/en/latest/api/vllm_omni/diffusion/worker/request_batch/#vllm_omni.diffusion.worker.request_batch.DiffusionRequestBatch) exposes compatibility properties such as `prompts`, `sampling_params`, and `request_id`. Pipelines that can execute the whole request batch in one forward pass should set `supports_request_batch = True`; other pipelines still receive a single-request batch and return a one-element output list.

See some parameters in `OmniDiffusionSamplingParams` as follows:

Expand All @@ -367,19 +369,18 @@ See some parameters in `OmniDiffusionSamplingParams` as follows:
**Extract parameters from request:**

```python
from vllm_omni.diffusion.request import OmniDiffusionRequest
from vllm_omni.diffusion.data import DiffusionOutput
from vllm_omni.diffusion.worker.request_batch import DiffusionRequestBatch

def forward(
self,
req: OmniDiffusionRequest,
) -> DiffusionOutput:
# Extract prompts from request
if req.prompts is not None:
prompt = [
p if isinstance(p, str) else (p.get("prompt") or "")
for p in req.prompts
]
req: DiffusionRequestBatch,
) -> list[DiffusionOutput]:
# Extract prompts from the request batch
prompts = [
p if isinstance(p, str) else (p.get("prompt") or "")
for p in req.prompts
]

# Extract sampling parameters
sampling_params = req.sampling_params
Expand All @@ -388,14 +389,16 @@ def forward(
height = sampling_params.height or (self.default_sample_size * self.vae_scale_factor)
width = sampling_params.width or (self.default_sample_size * self.vae_scale_factor)

# For image editing pipelines, extract images from multi_modal_data
if hasattr(req, 'multi_modal_data') and req.multi_modal_data:
input_images = req.multi_modal_data.get('image', [])
# For image editing pipelines, extract media from each prompt dict
input_images = []
for p in req.prompts:
multi_modal_data = p.get("multi_modal_data", {}) if isinstance(p, dict) else {}
input_images.append(multi_modal_data.get("image"))

# ... rest of generation logic
```

For an image editing model, an example `OmniDiffusionRequest` is like:
For an image editing model, the request `prompt` can be a dict like:
```python
{
"prompt": "turn this cat to a dog",
Expand Down Expand Up @@ -472,12 +475,12 @@ def get_your_model_pre_process_func(
def pre_process_func(
request: OmniDiffusionRequest,
):
for i, prompt in enumerate(request.prompts):
multi_modal_data = prompt.get("multi_modal_data", {}) if not isinstance(prompt, str) else None
raw_image = multi_modal_data.get("image", None) if multi_modal_data is not None else None
# image pre-processing
# after pre-processing, update the request attributes
...
prompt = request.prompt
multi_modal_data = prompt.get("multi_modal_data", {}) if not isinstance(prompt, str) else None
raw_image = multi_modal_data.get("image", None) if multi_modal_data is not None else None
# image pre-processing
# after pre-processing, update the request attributes
...
return request

return pre_process_func
Expand Down Expand Up @@ -923,11 +926,11 @@ When implementing a new pipeline, avoid putting all logic inside a single functi

For example:
```
def forward(self, req: OmniDiffusionRequest):
def forward(self, req: DiffusionRequestBatch) -> list[DiffusionOutput]:
prompt_embeds = self.encode_prompt(req)
latents = self.diffuse(prompt_embeds, req)
images = self.vae.decode(latents)
return DiffusionOutput(output=images)
return [DiffusionOutput(output=images)]
```
This allows the timing utility to measure each stage (e.g., encode_prompt, diffuse, vae.decode) separately and helps identify performance bottlenecks more easily.

Expand Down
22 changes: 13 additions & 9 deletions docs/design/feature/diffusion_continuous_batching.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,9 +38,13 @@ With continuous batching enabled:
The current implementation is conservative:

- only compatible requests are batched together
- request-mode diffusion still runs with `max_num_seqs=1`
- per-request progress and completion remain independent

Here, "continuous batching" means the step-wise path enabled by
`step_execution=True`. Request-mode `DiffusionRequestBatch` is static
request-level batching for one full pipeline `forward()` call; it does not
admit or remove requests between denoise steps.

## Enablement

Use `--step-execution` as the feature gate, then increase `--max-num-seqs`
Expand Down Expand Up @@ -72,7 +76,7 @@ which is built from shape-sensitive and CFG-sensitive sampling fields. This is
the core correctness rule for batching: requests are only co-batched when they
share the same denoise tensor contract.

There are two important details:
There are three important details:

- `num_inference_steps` is not part of the key, so requests with different
total step counts can still share a batch
Expand All @@ -88,8 +92,10 @@ key also covers LoRA identity (`lora_int_id`, `lora_scale`), so requests
targeting different adapters or scales run in separate batches and the
worker can activate exactly one adapter per step.

The current batching unit is one `OmniDiffusionRequest`. Requests with
multiple prompts do not participate in batching today.
The scheduler batching unit is one logical `OmniDiffusionRequest`. In the
step-wise path, runtime tensor batching is represented as `StepInputBatch`. For
request-mode prompt semantics, see
[Request-Level Batching](../../user_guide/diffusion/request_batching.md).

## Runner

Expand Down Expand Up @@ -126,17 +132,15 @@ request-local scheduler state and outputs.
the background loop and async add-request path needed for multiple requests to
accumulate in the scheduler.

This is supporting infrastructure, not the main design point. The batching
behavior is defined by scheduler-side compatibility gating and runner-side
batch packing.
When `step_execution=True`, the engine routes work through the step-wise
executor path. The continuous batching behavior is defined by scheduler-side
compatibility gating and runner-side `StepInputBatch` packing.

## Current Limitations

- Experimental feature; use `max_num_seqs=1` for the older conservative path.
- Only native pipelines that already support `step_execution=True`.
- Request-mode diffusion still clamps `max_num_seqs` back to `1`.
- Only homogeneous batches keyed by `SamplingParamsKey` are supported.
- Multi-prompt requests are not batched.
- `cache_backend`, KV transfer, and other request-mode extras are not wired
into the batched step-wise path yet.
- Future work can relax the current same-shape restriction with richer
Expand Down
175 changes: 175 additions & 0 deletions docs/design/feature/diffusion_request_level_batching.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,175 @@
# Request-Level Batching for Diffusion

This document describes the request-mode batching path for diffusion pipelines.
For end-user enablement and tuning, see
[Request-Level Batching](../../user_guide/diffusion/request_batching.md).

This is separate from
[Continuous Batching for Step-Wise Diffusion](diffusion_continuous_batching.md).
Request-level batching runs one full pipeline `forward()` over a static batch of
compatible requests. Step-wise continuous batching admits work between denoise
steps when `step_execution=True`.

## Why It Helps

The request-level design avoids coupling several logical requests to one request
object. This keeps request identity, abort/error handling, and per-request
metadata unambiguous while still allowing one fused pipeline forward pass for
bursty or concurrent traffic.

## Overview

With request-level batching enabled:

- each `OmniDiffusionRequest` contains one `prompt` and one `request_id`
- the scheduler groups compatible waiting requests into one scheduler wave
- `DiffusionRequestBatch` wraps the scheduled requests for pipeline `forward()`
- batch-capable pipelines return `list[DiffusionOutput]`, one output per
request
- `BatchRunnerOutput` maps each result back to its original `request_id`

Pipelines opt in with `supports_request_batch = True` and a `forward()` method
that accepts `DiffusionRequestBatch` and returns `list[DiffusionOutput]`.
Pipelines that do not opt in keep the existing per-request execution path.

## Enablement

Request-level batching is the request-mode path, so `step_execution` must remain
disabled. Increase `max_num_seqs` above `1` to let the scheduler keep multiple
compatible requests active:

```bash
vllm serve Qwen/Qwen-Image --omni \
--port 8091 \
--max-num-seqs 4
```

For bursty online ingress, `request_batch_max_wait_ms` can add a bounded
admission wait before the first `schedule()` of a scheduler wave:

```bash
vllm serve Qwen/Qwen-Image --omni \
--port 8091 \
--max-num-seqs 4 \
--request-batch-max-wait-ms 20
```

`request_batch_max_wait_ms=0` disables this wait and is the default.

## Request Contract

`OmniDiffusionRequest` represents one logical request. It owns one prompt,
sampling parameters, request id, and request-local metadata. Runtime batches are
formed by the scheduler and represented separately from the request payload.

Runtime batching is represented by:

- [`DiffusionSchedulerOutput`](gh-file:vllm_omni/diffusion/sched/interface.py)
for scheduled request ids and request payloads
- [`DiffusionRequestBatch`](gh-file:vllm_omni/diffusion/worker/request_batch.py)
for the pipeline-facing request batch
- [`BatchRunnerOutput`](gh-file:vllm_omni/diffusion/worker/utils.py) for
per-request results

`DiffusionRequestBatch` intentionally exposes compatibility properties such as
`prompts`, `sampling_params`, `request_id`, and `kv_sender_info` so migrated
pipelines can stay close to upstream code while using a batch-aware contract.

## Scheduler

The scheduler derives its capacity from `max_num_seqs` through
`max_num_running_reqs`. It exposes waiting/running queue counters so the engine
can decide whether admission wait is useful before scheduling a new wave.

Batch compatibility is controlled by
[`SamplingParamsKey`](gh-file:vllm_omni/diffusion/sched/interface.py). The key
contains shape-sensitive and guidance-sensitive fields, including output count
and LoRA identity. Requests with incompatible shapes, CFG settings, output
counts, LoRA adapters, or LoRA scales are kept in separate batches.

Admission is conservative:

- the scheduler only batches compatible requests
- FIFO ordering is preserved
- an incompatible request at the head of the waiting queue blocks later
compatible requests

## Engine

[`DiffusionEngine`](gh-file:vllm_omni/diffusion/diffusion_engine.py) resolves
request-batch capability during initialization from the configured pipeline
class, including custom pipeline classes.

The capability check uses the pipeline class attribute
`supports_request_batch = True`. Pipelines that set this attribute must implement
a request-batch-compatible `forward()` contract and return one
`DiffusionOutput` per request; the runner validates that return shape at runtime.

When the selected pipeline is batch-capable and `step_execution=False`, request
mode routes scheduler waves through the batch executor path. Otherwise it keeps
the per-request executor path.

The optional admission wait runs only when:

- request batching is supported
- `step_execution=False`
- `request_batch_max_wait_ms > 0`
- no requests are currently running

The wait exits early when the waiting queue reaches capacity, when the queue is
stable for a short window, when the deadline expires, or when the engine stops.

## Executor And Runner

The executor exposes two request-mode entries:

- `execute_request`: one worker call per scheduled request
- `execute_batch`: one worker call for the whole `DiffusionSchedulerOutput`

On the batch path, the worker builds a `DiffusionRequestBatch` and runs the
pipeline once. Request-local setup remains per request:

- KV transfer metadata
- random generator and seed handling
- request output/error/abort mapping

Shared batch setup happens once per batch when possible:

- cache refresh
- LoRA activation for the homogeneous adapter key
- pipeline `forward(req_batch)`

Large tensor IPC still uses the shared-memory packing path. The packer traverses
both normal `RunnerOutput.result` wrappers and nested batch results so batched
outputs do not fall back to pickle IPC for tensor payloads.

## Current Limitations

- Only pipelines that declare the request-batch contract use fused batch
execution.
- Batches are homogeneous under `SamplingParamsKey`; heterogeneous resolution or
incompatible guidance settings do not co-batch yet.
- FIFO scheduling can reduce batching opportunities when an incompatible
request is at the front of the queue.
- `request_batch_max_wait_ms` improves burst coalescing but can add latency to
the first request in a scheduler wave. Keep it small for latency-sensitive
serving.
- Step-wise continuous batching is documented separately and only applies when
`step_execution=True`.

## Related Files

- Request object and request batch:
[`vllm_omni/diffusion/request.py`](gh-file:vllm_omni/diffusion/request.py)
- Scheduler interface:
[`vllm_omni/diffusion/sched/interface.py`](gh-file:vllm_omni/diffusion/sched/interface.py)
- Scheduler base:
[`vllm_omni/diffusion/sched/base_scheduler.py`](gh-file:vllm_omni/diffusion/sched/base_scheduler.py)
- Engine:
[`vllm_omni/diffusion/diffusion_engine.py`](gh-file:vllm_omni/diffusion/diffusion_engine.py)
- Worker runner:
[`vllm_omni/diffusion/worker/diffusion_model_runner.py`](gh-file:vllm_omni/diffusion/worker/diffusion_model_runner.py)
- Executor interface:
[`vllm_omni/diffusion/executor/abstract.py`](gh-file:vllm_omni/diffusion/executor/abstract.py)
- Tests:
[`tests/diffusion/test_diffusion_engine.py`](gh-file:tests/diffusion/test_diffusion_engine.py)
1 change: 1 addition & 0 deletions docs/design/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ This section contains design documents and architecture specifications for vLLM-
- [Disaggregated Inference](feature/disaggregated_inference.md)
- [Ray-based Execution](feature/ray_based_execution.md)
- [Adding Step Execution Support for Diffusion Pipelines](feature/diffusion_step_execution.md)
- [Request-Level Batching for Diffusion](feature/diffusion_request_level_batching.md)
- [Continuous Batching for Step-Wise Diffusion](feature/diffusion_continuous_batching.md)

## Infrastructure Design Documents
Expand Down
Loading
Loading