diff --git a/docs/cli/bench/serve.md b/docs/cli/bench/serve.md index b9c59009587..a162c001a0f 100644 --- a/docs/cli/bench/serve.md +++ b/docs/cli/bench/serve.md @@ -436,6 +436,76 @@ This should be seen as an edge case, and if this behavior can be avoided by sett +### Image / video generation + +For `/v1/images/generations`, `/v1/images/edits`, and `/v1/videos`, set `--endpoint` to that path and omit `--backend`. The endpoint path is registered as the bench request adapter, so a separate named backend is not required. + +Example (`/v1/images/generations`): + +```bash +vllm bench serve --omni \ + --endpoint /v1/images/generations \ + --dataset-name random \ + --model ~/models/Qwen/Qwen-Image \ + --tokenizer ~/models/Qwen/Qwen-Image/tokenizer \ + --max-concurrency 1 \ + --num-warmups 2 \ + --num-prompts 10 \ + --random-input-len 64 \ + --random-output-len 1 \ + --ignore-eos \ + --percentile-metrics e2el \ + --extra-body '{ + "num_inference_steps": 20, + "seed": 42, + "true_cfg_scale": 4.0 + }' +``` + +If successful, following output is like: + +```text +============ Serving Benchmark Result ============ +Successful requests: 3 +Failed requests: 0 +Maximum request concurrency: 1 +Benchmark duration (s): 8.42 +Request throughput (req/s): 0.36 +Peak concurrent requests: 2.00 +-------------------Peak Memory-------------------- +Mean PEAK_MEMORY_MB (MB): 58832.00 +Median PEAK_MEMORY_MB (MB): 58832.00 +P99 PEAK_MEMORY_MB (MB): 58832.00 +----------------End-to-end Latency---------------- +Mean E2EL (ms): 2798.16 +Median E2EL (ms): 2799.83 +P99 E2EL (ms): 2800.03 +================== Image Result ================== +Total images generated: 3 +Image throughput (img/s): 0.36 +Average pixels per image: 1048576.00 +Mean denoise step latency (ms): 136.59 +---------------- Image Generation ---------------- +Mean IMAGE_GENERATION (ms): 2731.84 +Median IMAGE_GENERATION (ms): 2732.02 +P99 IMAGE_GENERATION (ms): 2734.65 +================================================== +``` + +Use the same pattern with `--endpoint /v1/images/edits` or `--endpoint /v1/videos` (and model-specific `--extra-body` as needed). Pure image/video runs omit the Text Result section when there is no generated text. + +`/v1/images/edits` defaults to **non-streaming JSON** (`stream=false`) so single-stage edit models are not rejected by the server. For multi-stage pipelines that need SSE (AR TTFT / image chunks), pass `"stream": true` in `--extra-body`. + +`/v1/videos` is an async job API: the client creates a job, then polls until `completed`/`failed`. Measured **e2el** therefore includes client poll sleep and any overshoot after the job actually finishes. Tune polling via `--extra-body`: + +- `poll_interval_s` (default `2.0`): sleep between status polls +- `poll_timeout_s` (default `21600`, i.e. 6 hours): give up waiting for the job + +**VIDEO_RTF** prefers server-reported generation time when available: + +- `video_rtf = video_generation_time_ms / 1000 / video_duration`, where `video_generation_time_ms` comes from response `stage_durations` (or `inference_time_s` as fallback) +- If generation time is missing, falls back to `e2el / video_duration` (this fallback *does* include poll overhead) + ### Multi-Stage Benchmark
diff --git a/tests/benchmarks/metrics/test_metrics.py b/tests/benchmarks/metrics/test_metrics.py index 31568525867..d7dd6c14a6b 100644 --- a/tests/benchmarks/metrics/test_metrics.py +++ b/tests/benchmarks/metrics/test_metrics.py @@ -457,5 +457,136 @@ def test_text_benchmark_still_reports_ttft(capsys): assert "Time to First Token" in out, "text benchmarks must keep TTFT" +# ============================================================================ +# Suppress text metrics for pure image / video endpoints +# ============================================================================ + + +def _make_image_output(prompt_len: int) -> MixRequestFuncOutput: + output = MixRequestFuncOutput() + output.success = True + output.prompt_len = prompt_len + output.output_tokens = 0 + output.generated_text = "" + output.ttft = 0.0 + output.text_latency = 0.0 + output.latency = 2.5 + output.start_time = 0.0 + output.itl = [] + output.image_count = 1 + output.image_generation_time_ms = 2000.0 + output.image_pixels = 1024 * 1024 + output.input_audio_duration = 0.0 + output.error = "" + return output + + +def _make_video_output(prompt_len: int) -> MixRequestFuncOutput: + output = MixRequestFuncOutput() + output.success = True + output.prompt_len = prompt_len + output.output_tokens = 0 + output.generated_text = "" + output.ttft = 0.0 + output.text_latency = 0.0 + output.latency = 5.0 + output.start_time = 0.0 + output.itl = [] + output.video_duration = 2.0 + output.video_frames = 48 + output.video_generation_time_ms = 4000.0 + output.video_rtf = 2.0 + output.input_audio_duration = 0.0 + output.error = "" + return output + + +_MEDIA_PERCENTILE_METRICS = ["e2el", "ttft"] + + +def test_pure_image_benchmark_omits_text_result(capsys): + """Pure image generation must not print Text Result or invent peak tok/s.""" + outputs = [_make_image_output(64), _make_image_output(64)] + + metrics, actual_output_lens = calculate_metrics( + input_requests=[], + outputs=outputs, + dur_s=5.0, + tokenizer=None, + selected_percentiles=[99.0], + goodput_config_dict={}, + task_type=TaskType.GENERATION, + selected_percentile_metrics=_MEDIA_PERCENTILE_METRICS, + max_concurrency=None, + request_rate=float("inf"), + benchmark_duration=5.0, + ) + + out = capsys.readouterr().out + assert actual_output_lens == [0, 0] + assert metrics.total_output == 0 + assert metrics.max_output_tokens_per_s == 0.0 + assert " Text Result " not in out + assert "Peak output token throughput" not in out + assert "Time to First Token" not in out + assert " Image Result " in out + + +def test_pure_video_benchmark_omits_text_result(capsys): + """Pure video generation must not print Text Result or invent peak tok/s.""" + outputs = [_make_video_output(63), _make_video_output(63)] + + metrics, actual_output_lens = calculate_metrics( + input_requests=[], + outputs=outputs, + dur_s=10.0, + tokenizer=None, + selected_percentiles=[99.0], + goodput_config_dict={}, + task_type=TaskType.GENERATION, + selected_percentile_metrics=_MEDIA_PERCENTILE_METRICS, + max_concurrency=None, + request_rate=float("inf"), + benchmark_duration=10.0, + ) + + out = capsys.readouterr().out + assert actual_output_lens == [0, 0] + assert metrics.total_output == 0 + assert metrics.max_output_tokens_per_s == 0.0 + assert " Text Result " not in out + assert "Peak output token throughput" not in out + assert "Time to First Token" not in out + assert " Video Result " in out + + +def test_image_with_generated_text_still_reports_text_result(capsys): + """Image edits / AR text alongside images must keep Text Result.""" + output = _make_image_output(64) + output.output_tokens = 8 + output.generated_text = "revised caption" + output.ttft = 0.05 + output.itl = [0.01] * 7 + + calculate_metrics( + input_requests=[], + outputs=[output], + dur_s=3.0, + tokenizer=_EmptyAwareTokenizer(), + selected_percentiles=[99.0], + goodput_config_dict={}, + task_type=TaskType.GENERATION, + selected_percentile_metrics=_MEDIA_PERCENTILE_METRICS, + max_concurrency=None, + request_rate=float("inf"), + benchmark_duration=3.0, + ) + + out = capsys.readouterr().out + assert " Text Result " in out + assert "Time to First Token" in out + assert " Image Result " in out + + if __name__ == "__main__": pytest.main([__file__, "-v", "-s"]) diff --git a/tests/benchmarks/patch/test_patch.py b/tests/benchmarks/patch/test_patch.py index e72d5a26347..b4582913ea1 100644 --- a/tests/benchmarks/patch/test_patch.py +++ b/tests/benchmarks/patch/test_patch.py @@ -22,9 +22,14 @@ ) from vllm_omni.benchmarks.patch.patch import ( MixRequestFuncOutput, + _add_video_extra_body_to_form, + _add_video_reference_to_form, _apply_stage0_token_timings, + _apply_video_metrics_from_payload, _attach_seed_tts_to_request_func_input, async_request_openai_chat_omni_completions, + async_request_openai_image_edits_omni, + async_request_openai_image_generations_omni, async_request_openai_realtime_duplex, should_request_stage_metrics, ) @@ -1196,5 +1201,272 @@ async def test_prompt_len_assigned_from_usage(mocker: MockerFixture): ) +def test_video_rtf_prefers_generation_time_over_poll_latency(): + """RTF must not grow with client poll_interval overshoot in output.latency.""" + payload = { + "duration_s": 2.0, + "num_frames": 48, + "fps": 24.0, + "stage_durations": {"stage_0_gen_ms": 4000.0}, + } + request_body: dict[str, object] = {} + + short_poll = MixRequestFuncOutput() + short_poll.latency = 4.2 # ~generation + small poll overshoot + _apply_video_metrics_from_payload(short_poll, payload, request_body) + + long_poll = MixRequestFuncOutput() + long_poll.latency = 8.0 # same job, larger poll_interval_s overshoot + _apply_video_metrics_from_payload(long_poll, payload, request_body) + + assert short_poll.video_generation_time_ms == pytest.approx(4000.0) + assert long_poll.video_generation_time_ms == pytest.approx(4000.0) + assert short_poll.video_rtf == pytest.approx(2.0) + assert long_poll.video_rtf == pytest.approx(2.0) + assert short_poll.video_rtf == long_poll.video_rtf + + +def test_video_rtf_falls_back_to_e2e_latency_without_generation_time(): + output = MixRequestFuncOutput() + output.latency = 6.0 + _apply_video_metrics_from_payload( + output, + {"duration_s": 2.0, "num_frames": 48, "fps": 24.0}, + {}, + ) + assert output.video_generation_time_ms == 0.0 + assert output.video_rtf == pytest.approx(3.0) + + +# 1x1 PNG used as a valid b64_json image payload in edit-client tests. +_MIN_PNG_B64 = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mP8z8BQDwAEhQGAhKmMIQAAAABJRU5ErkJggg==" + + +@pytest.mark.asyncio +async def test_image_edits_defaults_to_non_streaming_json(mocker: MockerFixture) -> None: + """Single-stage servers reject stream=true; default path must send stream=false + JSON.""" + + class MockJsonResponse: + status = 200 + + def __init__(self): + self._payload = { + "created": 1, + "data": [ + { + "b64_json": _MIN_PNG_B64, + "stage_durations": {"stage_0_gen_ms": 12.5}, + } + ], + } + + async def json(self): + return self._payload + + async def text(self): + return json.dumps(self._payload) + + async def __aenter__(self): + return self + + async def __aexit__(self, *_args): + return None + + captured_stream: list[str] = [] + real_add_field = None + + def tracking_add_field(self, name, value=None, **kwargs): + if name == "stream": + captured_stream.append(str(value)) + return real_add_field(self, name, value, **kwargs) + + import aiohttp + + real_add_field = aiohttp.FormData.add_field + mocker.patch.object(aiohttp.FormData, "add_field", tracking_add_field) + + mock_session = mocker.AsyncMock() + mock_session.post = mocker.MagicMock(return_value=MockJsonResponse()) + + request = RequestFuncInput( + model="single-stage-edit", + model_name="single-stage-edit", + prompt="make it sunny", + api_url="http://test.com/v1/images/edits", + prompt_len=4, + output_len=1, + multi_modal_content=[{"type": "image_url", "image_url": {"url": f"data:image/png;base64,{_MIN_PNG_B64}"}}], + extra_body={"num_inference_steps": 4}, + ) + output = await async_request_openai_image_edits_omni(request, mock_session, pbar=None) + + assert captured_stream == ["false"] + assert output.success is True + assert not output.error + assert output.image_count == 1 + assert output.image_generation_time_ms == pytest.approx(12.5) + + +@pytest.mark.asyncio +async def test_image_edits_stream_true_uses_sse_path(mocker: MockerFixture) -> None: + """Explicit stream=true keeps the multi-stage SSE client path.""" + sse_chunk = ( + b'data: {"type":"image","data":[{"b64_json":"' + _MIN_PNG_B64.encode() + b'"}]}\n\n' + b"data: [DONE]\n\n" + ) + mock_response = MockResponse(200, [sse_chunk]) + captured_stream: list[str] = [] + + import aiohttp + + real_add_field = aiohttp.FormData.add_field + + def tracking_add_field(self, name, value=None, **kwargs): + if name == "stream": + captured_stream.append(str(value)) + return real_add_field(self, name, value, **kwargs) + + mocker.patch.object(aiohttp.FormData, "add_field", tracking_add_field) + + mock_session = mocker.AsyncMock() + mock_session.post = mocker.MagicMock(return_value=mock_response) + + request = RequestFuncInput( + model="multi-stage-edit", + model_name="multi-stage-edit", + prompt="edit", + api_url="http://test.com/v1/images/edits", + prompt_len=2, + output_len=1, + multi_modal_content=[{"type": "image_url", "image_url": {"url": f"data:image/png;base64,{_MIN_PNG_B64}"}}], + extra_body={"stream": True}, + ) + output = await async_request_openai_image_edits_omni(request, mock_session, pbar=None) + + assert captured_stream == ["true"] + assert output.success is True + assert output.image_count == 1 + + +@pytest.mark.asyncio +async def test_image_generations_e2el_includes_json_body_consume(mocker: MockerFixture) -> None: + """E2EL must include body transfer/decode, not stop at HTTP headers.""" + + class SlowJsonResponse: + status = 200 + + def __init__(self): + self._payload = { + "created": 1, + "data": [{"b64_json": _MIN_PNG_B64, "stage_durations": {"stage_0_gen_ms": 1.0}}], + } + + async def json(self): + await asyncio.sleep(0.05) + return self._payload + + async def text(self): + return json.dumps(self._payload) + + async def __aenter__(self): + return self + + async def __aexit__(self, *_args): + return None + + mock_session = mocker.AsyncMock() + mock_session.post = mocker.MagicMock(return_value=SlowJsonResponse()) + request = RequestFuncInput( + model="img-gen", + model_name="img-gen", + prompt="a cat", + api_url="http://test.com/v1/images/generations", + prompt_len=2, + output_len=1, + extra_body={}, + ) + output = await async_request_openai_image_generations_omni(request, mock_session, pbar=None) + assert output.success is True + assert output.latency >= 0.05 + + +def test_video_local_image_reference_not_forwarded_as_raw_extra_field(tmp_path, mocker: MockerFixture) -> None: + """Local image_reference must only be uploaded via dedicated serializer.""" + import aiohttp + + ref_path = tmp_path / "ref.png" + ref_path.write_bytes(base64.b64decode(_MIN_PNG_B64)) + + captured: list[tuple[str, object]] = [] + real_add_field = aiohttp.FormData.add_field + + def tracking_add_field(self, name, value=None, **kwargs): + captured.append((str(name), value)) + return real_add_field(self, name, value, **kwargs) + + mocker.patch.object(aiohttp.FormData, "add_field", tracking_add_field) + + form = aiohttp.FormData() + extra_body = { + "num_inference_steps": 2, + "image_reference": str(ref_path), + } + request_body = {"model": "vid", "prompt": "p", **extra_body} + _add_video_extra_body_to_form(form, extra_body, request_body) + assert _add_video_reference_to_form(form, extra_body["image_reference"]) is True + + field_names = [name for name, _ in captured] + assert "image_reference" not in field_names + assert field_names.count("input_reference") == 1 + # Dedicated uploader sends file bytes, not the raw local path string. + uploaded = next(value for name, value in captured if name == "input_reference") + assert uploaded == base64.b64decode(_MIN_PNG_B64) + + +@pytest.mark.parametrize( + "reference", + [ + {"image_url": "https://example.com/ref.png"}, + {"file_id": "file-abc"}, + [{"image_url": "https://example.com/a.png"}, {"file_id": "file-xyz"}], + ], +) +def test_video_structured_image_reference_serialized_to_form(reference: object, mocker: MockerFixture) -> None: + """Object-form image_reference must be JSON-serialized, not silently dropped.""" + import aiohttp + + captured: list[tuple[str, object]] = [] + real_add_field = aiohttp.FormData.add_field + + def tracking_add_field(self, name, value=None, **kwargs): + captured.append((str(name), value)) + return real_add_field(self, name, value, **kwargs) + + mocker.patch.object(aiohttp.FormData, "add_field", tracking_add_field) + + form = aiohttp.FormData() + extra_body = {"image_reference": reference} + request_body = {"model": "vid", "prompt": "p", **extra_body} + _add_video_extra_body_to_form(form, extra_body, request_body) + assert _add_video_reference_to_form(form, reference) is True + + field_names = [name for name, _ in captured] + # Reserved key must not be double-forwarded as a generic extra field; + # exactly one dedicated image_reference field. + assert field_names.count("image_reference") == 1 + assert "input_reference" not in field_names + payload = next(value for name, value in captured if name == "image_reference") + assert json.loads(payload) == reference + + +def test_video_unsupported_image_reference_raises() -> None: + import aiohttp + + form = aiohttp.FormData() + with pytest.raises(ValueError, match="Unsupported image_reference"): + _add_video_reference_to_form(form, {"not_a_supported_key": "x"}) + with pytest.raises(ValueError, match="Unsupported image_reference"): + _add_video_reference_to_form(form, "/tmp/does-not-exist-ref.png") + + if __name__ == "__main__": pytest.main([__file__, "-v", "-s"]) diff --git a/tests/benchmarks/test_serve_cli.py b/tests/benchmarks/test_serve_cli.py index 5365d56ac7d..b9c72b1bf32 100644 --- a/tests/benchmarks/test_serve_cli.py +++ b/tests/benchmarks/test_serve_cli.py @@ -167,6 +167,26 @@ def test_preprocess_serve_args_applies_safe_omniinteract_prompt_default( {"bot_task": "recaption"}, {"backend", "print_stage", "bot_task"}, ), + ( + [ + "--endpoint", + "/v1/images/edits", + "--image-edits-bot-task", + "think", + ], + {"bot_task": "think"}, + {"endpoint", "bot_task"}, + ), + ( + [ + "--backend", + "/v1/images/edits", + "--image-edits-bot-task", + "recaption", + ], + {"bot_task": "recaption"}, + {"backend", "bot_task"}, + ), ( ["--extra-body", '{"bot_task":"vanilla"}'], {"bot_task": "vanilla"}, @@ -182,6 +202,7 @@ def test_omni_args_parse_and_preprocess( parser = TrackingArgumentParser() parser.add_argument("--extra-body", type=json.loads, default=None) parser.add_argument("--backend", default="openai-chat-omni") + parser.add_argument("--endpoint", default="/v1/chat/completions") add_omni_args(parser) args = parser.parse_args(argv) diff --git a/tests/entrypoints/openai_api/test_video_server.py b/tests/entrypoints/openai_api/test_video_server.py index 4a5931d6c09..b13dfe42e5d 100644 --- a/tests/entrypoints/openai_api/test_video_server.py +++ b/tests/entrypoints/openai_api/test_video_server.py @@ -1818,6 +1818,24 @@ def test_video_request_validation(): VideoGenerationRequest(prompt="test", quality="medium") +def test_async_create_accepts_fractional_fps(test_client, mocker: MockerFixture): + """Queued VideoResponse must accept fractional fps from the request path.""" + _mock_encode_video_bytes(mocker) + response = test_client.post( + "/v1/videos", + data={"prompt": "fractional fps", "fps": "12.5", "num_frames": "5"}, + ) + assert response.status_code == 200, response.text + body = response.json() + assert body["fps"] == 12.5 + assert body["num_frames"] == 5 + video_id = body["id"] + _wait_for_status(test_client, video_id, VideoGenerationStatus.COMPLETED.value) + engine = test_client.app.state.openai_serving_video._engine_client + assert engine.captured_sampling_params_list[0].fps == 12.5 + assert engine.captured_sampling_params_list[0].frame_rate == 12.5 + + def test_list_videos_supports_order_after_and_limit(test_client, mocker: MockerFixture): mocker.patch( "vllm_omni.entrypoints.openai.serving_video._encode_video_bytes", diff --git a/tests/metrics/test_metrics_utils.py b/tests/metrics/test_metrics_utils.py index 544c7819a7c..4ff201561ee 100644 --- a/tests/metrics/test_metrics_utils.py +++ b/tests/metrics/test_metrics_utils.py @@ -1,3 +1,6 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project + from __future__ import annotations from types import SimpleNamespace @@ -9,11 +12,30 @@ count_audio_frames, count_image_pixels, count_tokens_from_outputs, + count_video_frames, ) pytestmark = [pytest.mark.core_model, pytest.mark.cpu] +@pytest.mark.parametrize( + ("shape", "expected_frames"), + [ + ((1, 64, 64, 3), 1), + ((3, 64, 64, 3), 3), + ((4, 64, 64, 3), 4), + ((8, 64, 64, 3), 8), + ((3, 8, 64, 64), 8), + ((1, 1, 64, 64, 3), 1), + ((1, 3, 64, 64, 3), 3), + ((1, 4, 64, 64, 3), 4), + ((1, 3, 8, 64, 64), 8), + ], +) +def test_count_video_frames_handles_short_channel_last_outputs(shape, expected_frames): + assert count_video_frames(SimpleNamespace(shape=shape)) == expected_frames + + @pytest.mark.parametrize( ("audio_chunk", "expected"), [ diff --git a/vllm_omni/benchmarks/metrics/metrics.py b/vllm_omni/benchmarks/metrics/metrics.py index 3d7c983f3ea..52eec86b9c7 100644 --- a/vllm_omni/benchmarks/metrics/metrics.py +++ b/vllm_omni/benchmarks/metrics/metrics.py @@ -52,6 +52,21 @@ (defs.MEDIAN_IMAGE_GENERATION_MS, float, field(default=0.0)), (defs.STD_IMAGE_GENERATION_MS, float, field(default=0.0)), (defs.PERCENTILES_IMAGE_GENERATION_MS, _PERCENTILE_ROWS_TYPE, field(default=None)), + (defs.TOTAL_VIDEO_DURATION_S, float, field(default=0.0)), + (defs.TOTAL_VIDEO_FRAMES, int, field(default=0)), + (defs.VIDEO_THROUGHPUT, float, field(default=0.0)), + (defs.MEAN_VIDEO_RTF, float, field(default=0.0)), + (defs.MEDIAN_VIDEO_RTF, float, field(default=0.0)), + (defs.STD_VIDEO_RTF, float, field(default=0.0)), + (defs.PERCENTILES_VIDEO_RTF, _PERCENTILE_ROWS_TYPE, field(default=None)), + (defs.MEAN_VIDEO_GENERATION_MS, float, field(default=0.0)), + (defs.MEDIAN_VIDEO_GENERATION_MS, float, field(default=0.0)), + (defs.STD_VIDEO_GENERATION_MS, float, field(default=0.0)), + (defs.PERCENTILES_VIDEO_GENERATION_MS, _PERCENTILE_ROWS_TYPE, field(default=None)), + (defs.MEAN_PEAK_MEMORY_MB, float, field(default=0.0)), + (defs.MEDIAN_PEAK_MEMORY_MB, float, field(default=0.0)), + (defs.STD_PEAK_MEMORY_MB, float, field(default=0.0)), + (defs.PERCENTILES_PEAK_MEMORY_MB, _PERCENTILE_ROWS_TYPE, field(default=None)), ] MultiModalsBenchmarkMetrics = make_dataclass( @@ -171,8 +186,10 @@ def _stage_modality_flags( ) -> tuple[bool, bool, bool, bool, bool]: is_text_stage = final_output_type == "text" or output_unit_type == "text" is_audio_stage = final_output_type == "audio" or output_unit_type == "audio" - is_image_stage = final_output_type in {"image", "images"} or output_unit_type == "image" is_video_stage = final_output_type in {"video", "videos"} or output_unit_type == "video" + # Video diffusion may still report output_unit_type="image" when frames are + # stored in ``images``; prefer video when final_output_type says so. + is_image_stage = (not is_video_stage) and (final_output_type in {"image", "images"} or output_unit_type == "image") is_internal_stream_stage = ( output_unit_type in _STREAMING_OUTPUT_UNIT_TYPES and not is_text_stage and not is_audio_stage ) @@ -222,6 +239,7 @@ def print_metrics( print("{:<40} {:<10.2f}".format("Request goodput (req/s):", metrics.request_goodput)) if isinstance(metrics, MultiModalsBenchmarkMetrics): print("{:<40} {:<10.2f}".format("Peak concurrent requests:", metrics.max_concurrent_requests)) + print_peak_memory_metrics(metrics) if task_type != TaskType.GENERATION or "e2el" in selected_percentile_metrics: process_one_metric("e2el", metrics) print_text_metrics(task_type, selected_percentile_metrics, metrics) @@ -230,9 +248,13 @@ def print_metrics( print_audio_metrics(selected_percentile_metrics, metrics) if _has_image_output(metrics): print_image_metrics(selected_percentiles or [], metrics) + if _has_video_output(metrics): + print_video_metrics(selected_percentiles or [], metrics) if print_stage and outputs and selected_percentiles is not None: - print("\n{s:{c}^{n}}".format(s=" Stage Benchmark Result ", n=50, c="=")) - for sm in _build_stage_metrics_from_outputs(outputs): + stage_metrics = _build_stage_metrics_from_outputs(outputs) + if stage_metrics: + print("\n{s:{c}^{n}}".format(s=" Stage Benchmark Result ", n=50, c="=")) + for sm in stage_metrics: print_stage_metrics( task_type, selected_percentile_metrics, @@ -243,6 +265,11 @@ def print_metrics( def print_text_metrics(task_type, selected_percentile_metrics, metrics: MultiModalsBenchmarkMetrics): + # Pure image/video runs have no user-facing text tokens; skip the whole + # Text Result section (token throughput / peak would be misleading). + if metrics.total_output <= 0 and (_has_image_output(metrics) or _has_video_output(metrics)): + return + print("{s:{c}^{n}}".format(s=" Text Result ", n=50, c="=")) print("{:<40} {:<10}".format("Total input tokens:", metrics.total_input)) if isinstance(metrics, MultiModalsBenchmarkMetrics): @@ -275,6 +302,14 @@ def _has_image_output(metrics: MultiModalsBenchmarkMetrics) -> bool: return int(getattr(metrics, defs.TOTAL_IMAGES, 0) or 0) > 0 +def _has_video_output(metrics: MultiModalsBenchmarkMetrics) -> bool: + return ( + float(getattr(metrics, defs.TOTAL_VIDEO_DURATION_S, 0.0) or 0.0) > 0.0 + or int(getattr(metrics, defs.TOTAL_VIDEO_FRAMES, 0) or 0) > 0 + or float(getattr(metrics, defs.MEAN_VIDEO_GENERATION_MS, 0.0) or 0.0) > 0.0 + ) + + def print_audio_metrics(selected_percentile_metrics, metrics: MultiModalsBenchmarkMetrics): print("{s:{c}^{n}}".format(s=" Audio Result ", n=50, c="=")) print( @@ -293,6 +328,16 @@ def print_audio_metrics(selected_percentile_metrics, metrics: MultiModalsBenchma process_one_metric(metric, metrics) +def print_peak_memory_metrics(metrics: MultiModalsBenchmarkMetrics): + if getattr(metrics, defs.MEAN_PEAK_MEMORY_MB) <= 0: + return + print("{s:{c}^{n}}".format(s="Peak Memory", n=50, c="-")) + print("{:<40} {:<10.2f}".format("Mean PEAK_MEMORY_MB (MB):", getattr(metrics, defs.MEAN_PEAK_MEMORY_MB))) + print("{:<40} {:<10.2f}".format("Median PEAK_MEMORY_MB (MB):", getattr(metrics, defs.MEDIAN_PEAK_MEMORY_MB))) + for p, value in getattr(metrics, defs.PERCENTILES_PEAK_MEMORY_MB) or []: + print("{:<40} {:<10.2f}".format(f"P{_p_label(p)} PEAK_MEMORY_MB (MB):", value)) + + def print_image_metrics(selected_percentiles: list[float], metrics: MultiModalsBenchmarkMetrics): print("{s:{c}^{n}}".format(s=" Image Result ", n=50, c="=")) print("{:<40} {:<10}".format("Total images generated:", getattr(metrics, defs.TOTAL_IMAGES))) @@ -311,7 +356,7 @@ def print_image_metrics(selected_percentiles: list[float], metrics: MultiModalsB ) ) if getattr(metrics, defs.MEAN_IMAGE_GENERATION_MS) > 0: - print("-----------------Image Generation-----------------") + print("{s:{c}^{n}}".format(s=" Image Generation ", n=50, c="-")) print("{:<40} {:<10.2f}".format("Mean IMAGE_GENERATION (ms):", getattr(metrics, defs.MEAN_IMAGE_GENERATION_MS))) print( "{:<40} {:<10.2f}".format( @@ -324,6 +369,42 @@ def print_image_metrics(selected_percentiles: list[float], metrics: MultiModalsB print("{:<40} {:<10.2f}".format(f"P{_p_label(p)} IMAGE_GENERATION (ms):", value)) +def print_video_metrics(selected_percentiles: list[float], metrics: MultiModalsBenchmarkMetrics): + print("{s:{c}^{n}}".format(s=" Video Result ", n=50, c="=")) + print( + "{:<40} {:<10.2f}".format( + "Total video duration generated(s):", + getattr(metrics, defs.TOTAL_VIDEO_DURATION_S), + ) + ) + print("{:<40} {:<10}".format("Total video frames generated:", getattr(metrics, defs.TOTAL_VIDEO_FRAMES))) + print( + "{:<40} {:<10.2f}".format( + "Video throughput(video duration/s):", + getattr(metrics, defs.VIDEO_THROUGHPUT), + ) + ) + if getattr(metrics, defs.MEAN_VIDEO_RTF) > 0: + print("{s:{c}^{n}}".format(s=" Video RTF ", n=50, c="-")) + print("{:<40} {:<10.2f}".format("Mean VIDEO_RTF:", getattr(metrics, defs.MEAN_VIDEO_RTF))) + print("{:<40} {:<10.2f}".format("Median VIDEO_RTF:", getattr(metrics, defs.MEDIAN_VIDEO_RTF))) + for p, value in getattr(metrics, defs.PERCENTILES_VIDEO_RTF) or []: + if p in selected_percentiles: + print("{:<40} {:<10.2f}".format(f"P{_p_label(p)} VIDEO_RTF:", value)) + if getattr(metrics, defs.MEAN_VIDEO_GENERATION_MS) > 0: + print("{s:{c}^{n}}".format(s=" Video Generation ", n=50, c="-")) + print("{:<40} {:<10.2f}".format("Mean VIDEO_GENERATION (ms):", getattr(metrics, defs.MEAN_VIDEO_GENERATION_MS))) + print( + "{:<40} {:<10.2f}".format( + "Median VIDEO_GENERATION (ms):", + getattr(metrics, defs.MEDIAN_VIDEO_GENERATION_MS), + ) + ) + for p, value in getattr(metrics, defs.PERCENTILES_VIDEO_GENERATION_MS) or []: + if p in selected_percentiles: + print("{:<40} {:<10.2f}".format(f"P{_p_label(p)} VIDEO_GENERATION (ms):", value)) + + def process_one_metric( metric_attribute_name: str, metrics: MultiModalsBenchmarkMetrics, @@ -495,6 +576,21 @@ def _print_image_stage_metrics( ) +def _print_video_stage_metrics( + selected_percentiles: list[float], + sm: StageBenchmarkMetrics, +) -> None: + print("{s:{c}^{n}}".format(s=" Video Result ", n=50, c="=")) + _print_percentile_metric( + "Video Generation", + "VIDEO_GENERATION", + getattr(sm, defs.STAGE_GEN_TIMES_MS), + selected_percentiles, + values_are_ms=True, + to_ms=True, + ) + + def _print_internal_stream_stage_metrics( selected_percentile_metrics: list[str], selected_percentiles: list[float], @@ -549,11 +645,12 @@ def print_stage_metrics( ) = _stage_modality_flags(getattr(sm, "final_output_type"), getattr(sm, "output_unit_type")) print("{s:{c}^{n}}".format(s=title, n=50, c="=")) - if is_image_stage: - _print_image_stage_metrics(selected_percentiles, sm) + if is_video_stage: + _print_video_stage_metrics(selected_percentiles, sm) return - if is_video_stage: + if is_image_stage: + _print_image_stage_metrics(selected_percentiles, sm) return _print_stage_timing(sm, selected_percentiles) @@ -754,6 +851,11 @@ def calculate_metrics( denoise_step_latencies_ms: list[float] = [] total_images = 0 total_image_pixels = 0 + video_durations: list[float] = [] + video_rtfs: list[float] = [] + total_video_frames = 0 + video_generation_times_ms: list[float] = [] + peak_memories_mb: list[float] = [] audio_underruns: list[float] = [] audio_continuity_ok: list[bool] = [] input_audio_duration = 0.0 @@ -762,7 +864,18 @@ def calculate_metrics( output_len = outputs[i].output_tokens if not output_len: - if tokenizer is None: + # Pure image/video (no generated text) must stay at 0 tokens. + # Do not invent output_len=1 via the tokenizer-is-None placeholder. + generated_text = getattr(outputs[i], "generated_text", None) or "" + has_image = int(getattr(outputs[i], defs.IMAGE_COUNT, 0) or 0) > 0 + has_video = ( + float(getattr(outputs[i], defs.VIDEO_DURATION, 0.0) or 0.0) > 0.0 + or int(getattr(outputs[i], defs.VIDEO_FRAMES, 0) or 0) > 0 + or float(getattr(outputs[i], defs.VIDEO_GENERATION_TIME_MS, 0.0) or 0.0) > 0.0 + ) + if not generated_text and (has_image or has_video): + output_len = 0 + elif tokenizer is None: output_len = 1 else: # We use the tokenizer to count the number of output tokens @@ -827,6 +940,19 @@ def calculate_metrics( denoise_step_latency_ms = float(getattr(outputs[i], defs.DENOISE_STEP_LATENCY_MS, 0.0) or 0.0) if denoise_step_latency_ms > 0: denoise_step_latencies_ms.append(denoise_step_latency_ms) + video_duration = float(getattr(outputs[i], defs.VIDEO_DURATION, 0.0) or 0.0) + if video_duration > 0: + video_durations.append(video_duration) + total_video_frames += int(getattr(outputs[i], defs.VIDEO_FRAMES, 0) or 0) + video_rtf = float(getattr(outputs[i], defs.VIDEO_RTF, 0.0) or 0.0) + if video_rtf > 0: + video_rtfs.append(video_rtf) + video_generation_time_ms = float(getattr(outputs[i], defs.VIDEO_GENERATION_TIME_MS, 0.0) or 0.0) + if video_generation_time_ms > 0: + video_generation_times_ms.append(video_generation_time_ms) + peak_memory_mb = float(getattr(outputs[i], defs.PEAK_MEMORY_MB, 0.0) or 0.0) + if peak_memory_mb > 0: + peak_memories_mb.append(peak_memory_mb) audio_underruns.append(getattr(outputs[i], f"{defs.AUDIO_UNDERRUN}_s", 0.0)) audio_continuity_ok.append(bool(getattr(outputs[i], defs.AUDIO_CONTINUITY_OK, True))) e2els.append(outputs[i].latency) @@ -892,18 +1018,24 @@ def calculate_metrics( token_timeline_available = False else: # Calculate token generation timestamp using - # start_time, ttft, and itl - token_times = [output.start_time + output.ttft] - current_time = token_times[0] - for itl_value in output.itl: - current_time += itl_value - token_times.append(current_time) - - # Add tokens to second buckets - for token_time in token_times: - second_bucket = int(token_time - min_start_time) - if 0 <= second_bucket < duration_seconds: - tokens_per_second[second_bucket] += 1 + # start_time, ttft, and itl. Skip requests with no text tokens + # (empty itl, zero output_tokens, and no generated_text) so pure + # image/video runs do not seed a fake peak of 1 tok/s from + # start_time+ttft alone. + output_token_count = int(getattr(output, "output_tokens", 0) or 0) + has_generated_text = bool(getattr(output, "generated_text", None) or "") + if output_token_count > 0 or output.itl or has_generated_text: + token_times = [output.start_time + output.ttft] + current_time = token_times[0] + for itl_value in output.itl: + current_time += itl_value + token_times.append(current_time) + + # Add tokens to second buckets + for token_time in token_times: + second_bucket = int(token_time - min_start_time) + if 0 <= second_bucket < duration_seconds: + tokens_per_second[second_bucket] += 1 # Track concurrent requests for each second this request was active request_start_second = int(output.start_time - min_start_time) @@ -1009,6 +1141,25 @@ def calculate_metrics( defs.PERCENTILES_IMAGE_GENERATION_MS: [ (p, np.percentile(image_generation_times_ms or 0, p)) for p in selected_percentiles ], + defs.TOTAL_VIDEO_DURATION_S: sum(video_durations), + defs.TOTAL_VIDEO_FRAMES: total_video_frames, + defs.VIDEO_THROUGHPUT: sum(video_durations) / dur_s, + defs.MEAN_VIDEO_RTF: np.mean(video_rtfs or 0), + defs.STD_VIDEO_RTF: np.std(video_rtfs or 0), + defs.MEDIAN_VIDEO_RTF: np.median(video_rtfs or 0), + defs.PERCENTILES_VIDEO_RTF: [(p, np.percentile(video_rtfs or 0, p)) for p in selected_percentiles], + defs.MEAN_VIDEO_GENERATION_MS: np.mean(video_generation_times_ms or 0), + defs.STD_VIDEO_GENERATION_MS: np.std(video_generation_times_ms or 0), + defs.MEDIAN_VIDEO_GENERATION_MS: np.median(video_generation_times_ms or 0), + defs.PERCENTILES_VIDEO_GENERATION_MS: [ + (p, np.percentile(video_generation_times_ms or 0, p)) for p in selected_percentiles + ], + defs.MEAN_PEAK_MEMORY_MB: np.mean(peak_memories_mb or 0), + defs.STD_PEAK_MEMORY_MB: np.std(peak_memories_mb or 0), + defs.MEDIAN_PEAK_MEMORY_MB: np.median(peak_memories_mb or 0), + defs.PERCENTILES_PEAK_MEMORY_MB: [ + (p, np.percentile(peak_memories_mb or 0, p)) for p in selected_percentiles + ], defs.MEAN_AUDIO_UNDERRUN_S: np.mean(audio_underruns or 0), defs.STD_AUDIO_UNDERRUN_S: np.std(audio_underruns or 0), defs.MEDIAN_AUDIO_UNDERRUN_S: np.median(audio_underruns or 0), diff --git a/vllm_omni/benchmarks/patch/patch.py b/vllm_omni/benchmarks/patch/patch.py index b51dd2ee44f..e65e5ee5cfc 100644 --- a/vllm_omni/benchmarks/patch/patch.py +++ b/vllm_omni/benchmarks/patch/patch.py @@ -14,7 +14,7 @@ import traceback import uuid import wave -from collections.abc import Iterable +from collections.abc import Iterable, Mapping from dataclasses import dataclass, replace from datetime import datetime from pathlib import Path @@ -23,6 +23,7 @@ import aiohttp import numpy as np import pybase64 as base64 +from PIL import Image from tqdm.asyncio import tqdm from vllm.benchmarks import datasets from vllm.benchmarks.datasets import SampleRequest @@ -85,7 +86,7 @@ write_batch_artifacts as write_omniinteract_batch_artifacts, ) from vllm_omni.metrics import definitions as defs -from vllm_omni.metrics.utils import coerce_positive_int_scalar +from vllm_omni.metrics.utils import coerce_bool, coerce_positive_float_scalar, coerce_positive_int_scalar if TYPE_CHECKING: from vllm_omni.clients.duplex import DuplexClient @@ -100,7 +101,13 @@ _AUDIO_CONTINUITY_THRESHOLD_ENV = "VLLM_OMNI_BENCH_AUDIO_CONTINUITY_THRESHOLD_S" RETURN_STAGE_METRICS_FIELD = "return_stage_metrics" -_IMAGE_STAGE_METRICS_BACKENDS = frozenset({"openai-image-edits-omni"}) +_IMAGE_STAGE_METRICS_BACKENDS = frozenset( + { + "/v1/images/generations", + "/v1/images/edits", + "openai-image-edits-omni", + } +) _PRINT_STAGE = False @@ -662,6 +669,11 @@ class MixRequestFuncOutput(RequestFuncOutput): image_generation_time_ms: float = 0.0 image_pixels: int = 0 denoise_step_latency_ms: float = 0.0 + video_duration: float = 0.0 + video_rtf: float = 0.0 + video_frames: int = 0 + video_generation_time_ms: float = 0.0 + peak_memory_mb: float = 0.0 text_latency: float = 0.0 tpot_measured: bool = True #: Worst-case streaming-audio underrun (wall-clock seconds the player @@ -849,7 +861,7 @@ def _record_text_token_stream_intervals( def _update_output_stage_metrics_from_payload( output: MixRequestFuncOutput, - data: dict[str, Any], + data: Mapping[str, object], *, update_output_tokens: bool = True, ) -> None: @@ -897,7 +909,55 @@ def _apply_chat_stage0_token_timings(output: MixRequestFuncOutput) -> bool: ) -def _image_metrics_from_stage_metrics(metrics: dict[str, Any] | None) -> tuple[int, float, int, float]: +def _peak_memory_mb_from_payload(data: Mapping[str, object]) -> float: + peak_memory_mb = coerce_positive_float_scalar(data.get(defs.PEAK_MEMORY_MB)) + if peak_memory_mb is not None: + return peak_memory_mb + + for key in ("metrics", "usage"): + nested = data.get(key) + if isinstance(nested, dict): + peak_memory_mb = coerce_positive_float_scalar(nested.get(defs.PEAK_MEMORY_MB)) + if peak_memory_mb is not None: + return peak_memory_mb + + response_data = data.get("data") + if isinstance(response_data, list): + for item in response_data: + if isinstance(item, dict): + peak_memory_mb = coerce_positive_float_scalar(item.get(defs.PEAK_MEMORY_MB)) + if peak_memory_mb is not None: + return peak_memory_mb + + choices = data.get("choices") + if isinstance(choices, list): + for choice in choices: + if not isinstance(choice, dict): + continue + contents = [] + message = choice.get("message") + if isinstance(message, dict): + contents.append(message.get("content")) + delta = choice.get("delta") + if isinstance(delta, dict): + contents.append(delta.get("content")) + for content in contents: + if isinstance(content, list): + for item in content: + if isinstance(item, dict): + peak_memory_mb = coerce_positive_float_scalar(item.get(defs.PEAK_MEMORY_MB)) + if peak_memory_mb is not None: + return peak_memory_mb + return 0.0 + + +def _update_output_peak_memory_from_payload(output: MixRequestFuncOutput, data: Mapping[str, object]) -> None: + peak_memory_mb = _peak_memory_mb_from_payload(data) + if peak_memory_mb > output.peak_memory_mb: + output.peak_memory_mb = peak_memory_mb + + +def _image_metrics_from_stage_metrics(metrics: object) -> tuple[int, float, int, float]: if not isinstance(metrics, dict): return 0, 0.0, 0, 0.0 stage_snapshot = metrics.get("stage_metrics") @@ -924,7 +984,7 @@ def _image_metrics_from_stage_metrics(metrics: dict[str, Any] | None) -> tuple[i return image_count, image_generation_ms, image_pixels, denoise_step_latency_ms -def _image_generation_ms_from_content(content: Any) -> float: +def _image_generation_ms_from_content(content: object) -> float: if not isinstance(content, list): return 0.0 for item in content: @@ -943,6 +1003,263 @@ def _image_generation_ms_from_content(content: Any) -> float: return 0.0 +def _image_info_from_response_data(content: object) -> tuple[int, int]: + if not isinstance(content, list): + return 0, 0 + image_count = 0 + total_pixels = 0 + for item in content: + if not isinstance(item, dict): + continue + b64_json = item.get("b64_json") + if not isinstance(b64_json, str) or not b64_json: + continue + try: + with Image.open(io.BytesIO(base64.b64decode(b64_json, validate=True))) as img: + width, height = img.size + img.verify() + image_count += 1 + total_pixels += int(width) * int(height) + except Exception: + logger.debug("Failed to decode generated image payload", exc_info=True) + return image_count, total_pixels + + +def _apply_image_metrics_from_payload(output: MixRequestFuncOutput, data: Mapping[str, object]) -> int: + """Populate image benchmark fields from an OpenAI-compatible image payload.""" + _update_output_stage_metrics_from_payload(output, data, update_output_tokens=False) + _update_output_peak_memory_from_payload(output, data) + + payload_image_count = 0 + response_data = data.get("data") + if isinstance(response_data, list): + payload_image_count, content_image_pixels = _image_info_from_response_data(response_data) + output.image_count = max(output.image_count, payload_image_count) + content_image_ms = _image_generation_ms_from_content(response_data) + if content_image_ms > 0: + output.image_generation_time_ms = max(output.image_generation_time_ms, content_image_ms) + if content_image_pixels > 0: + output.image_pixels = max(output.image_pixels, content_image_pixels) + + ( + metrics_image_count, + metrics_image_ms, + metrics_image_pixels, + metrics_denoise_step_ms, + ) = _image_metrics_from_stage_metrics(data.get("metrics")) + if metrics_image_count > output.image_count: + output.image_count = metrics_image_count + if metrics_image_ms > output.image_generation_time_ms: + output.image_generation_time_ms = metrics_image_ms + if metrics_image_pixels > output.image_pixels: + output.image_pixels = metrics_image_pixels + if metrics_denoise_step_ms > output.denoise_step_latency_ms: + output.denoise_step_latency_ms = metrics_denoise_step_ms + return payload_image_count + + +_VIDEO_FORM_FIELDS = ( + "seconds", + "num_frames", + "fps", + "num_inference_steps", + "seed", + "negative_prompt", + "guidance_scale", + "guidance_scale_2", + "boundary_ratio", + "flow_shift", + "true_cfg_scale", + "generate_sound", + "sound_duration", + "enable_frame_interpolation", + "frame_interpolation_exp", + "frame_interpolation_scale", + "frame_interpolation_model_path", + "lora", + "extra_params", +) + + +def _video_generation_ms_from_stage_durations(stage_durations: object) -> float: + if not isinstance(stage_durations, dict): + return 0.0 + gen_values = [ + float(value) + for key, value in stage_durations.items() + if str(key).endswith("_gen_ms") and isinstance(value, (int, float)) + ] + return max(gen_values) if gen_values else 0.0 + + +def _video_duration_from_payload(data: Mapping[str, object], request_body: Mapping[str, object]) -> float: + duration_s = coerce_positive_float_scalar(data.get("duration_s")) + if duration_s is not None and duration_s > 0: + return duration_s + + num_frames = coerce_positive_float_scalar(data.get("num_frames")) + fps = coerce_positive_float_scalar(data.get("fps")) + if num_frames is not None and fps is not None and num_frames > 0 and fps > 0: + return num_frames / fps + + seconds = coerce_positive_float_scalar(request_body.get("seconds")) + if seconds is not None and seconds > 0: + return seconds + + num_frames = coerce_positive_float_scalar(request_body.get("num_frames")) + fps = coerce_positive_float_scalar(request_body.get("fps")) + if num_frames is not None and fps is not None and num_frames > 0 and fps > 0: + return num_frames / fps + + return 0.0 + + +def _video_frames_from_payload(data: Mapping[str, object], request_body: Mapping[str, object]) -> int: + for key in ("num_frames", "video_frames", "frames"): + value = data.get(key) + if isinstance(value, list): + return len(value) + num_frames = coerce_positive_int_scalar(value) + if num_frames is not None: + return num_frames + + duration_s = coerce_positive_float_scalar(data.get("duration_s")) + fps = coerce_positive_float_scalar(data.get("fps")) + if duration_s is not None and fps is not None and duration_s > 0 and fps > 0: + return int(round(duration_s * fps)) + + num_frames = coerce_positive_int_scalar(request_body.get("num_frames")) + if num_frames is not None: + return num_frames + + seconds = coerce_positive_float_scalar(request_body.get("seconds")) + fps = coerce_positive_float_scalar(request_body.get("fps")) + if seconds is not None and fps is not None and seconds > 0 and fps > 0: + return int(round(seconds * fps)) + + return 0 + + +def _is_structured_image_reference(reference: Mapping[str, object]) -> bool: + """True for API image_reference objects ({image_url}/{file_id}).""" + image_url = reference.get("image_url") + file_id = reference.get("file_id") + has_url = isinstance(image_url, str) and bool(image_url) + has_file_id = isinstance(file_id, str) and bool(file_id) + return has_url or has_file_id + + +def _add_video_reference_to_form(form: aiohttp.FormData, reference: object) -> bool: + if isinstance(reference, dict) and "bytes" in reference: + form.add_field( + "input_reference", + reference["bytes"], + filename="benchmark-reference", + content_type=reference.get("content_type", "application/octet-stream"), + ) + return True + + if isinstance(reference, Mapping) and _is_structured_image_reference(reference): + form.add_field("image_reference", json.dumps(dict(reference))) + return True + + if isinstance(reference, list): + if reference and all(isinstance(item, Mapping) and _is_structured_image_reference(item) for item in reference): + form.add_field("image_reference", json.dumps([dict(item) for item in reference])) + return True + raise ValueError( + "Unsupported image_reference list; expected non-empty list of " + '{"image_url": "..."} and/or {"file_id": "..."} objects.' + ) + + if isinstance(reference, str): + if reference.startswith(("data:image", "http://", "https://")): + form.add_field("image_reference", json.dumps({"image_url": reference})) + return True + local_path = reference.removeprefix("file://") + if os.path.exists(local_path): + with open(local_path, "rb") as f: + reference_bytes = f.read() + form.add_field( + "input_reference", + reference_bytes, + filename=os.path.basename(local_path), + content_type=_guess_mime_type(local_path), + ) + return True + raise ValueError(f"Unsupported image_reference path or URL: {reference!r}") + + raise ValueError( + "Unsupported image_reference; expected upload bytes, local path/URL string, " + 'or {"image_url": "..."} / {"file_id": "..."} object ' + f"(got {type(reference).__name__})." + ) + + +def _add_video_extra_body_to_form( + form: aiohttp.FormData, + extra_body: Mapping[str, object], + request_body: Mapping[str, object], +) -> None: + for key in _VIDEO_FORM_FIELDS: + value = request_body.get(key) + if value is None: + continue + if isinstance(value, (dict, list)): + form.add_field(key, json.dumps(value)) + else: + form.add_field(key, str(value)) + + reserved = { + "model", + "prompt", + "size", + "width", + "height", + "poll_interval_s", + "poll_timeout_s", + # Handled only by _add_video_reference_to_form (upload / JSON image_url). + "image_reference", + "input_reference", + *_VIDEO_FORM_FIELDS, + } + for key, value in extra_body.items(): + if key in reserved or value is None: + continue + if isinstance(value, (dict, list)): + form.add_field(key, json.dumps(value)) + else: + form.add_field(key, str(value)) + + +def _apply_video_metrics_from_payload( + output: MixRequestFuncOutput, + data: Mapping[str, object], + request_body: Mapping[str, object], +) -> None: + output.video_duration = _video_duration_from_payload(data, request_body) + output.video_frames = _video_frames_from_payload(data, request_body) + _update_output_stage_metrics_from_payload(output, data, update_output_tokens=False) + _update_output_peak_memory_from_payload(output, data) + + stage_durations = data.get("stage_durations") + stage_gen_ms = _video_generation_ms_from_stage_durations(stage_durations) + if stage_gen_ms <= 0: + inference_time_s = coerce_positive_float_scalar(data.get("inference_time_s")) + if inference_time_s is not None and inference_time_s > 0: + stage_gen_ms = inference_time_s * 1000.0 + output.video_generation_time_ms = max(output.video_generation_time_ms, stage_gen_ms) + if output.video_duration <= 0: + return + # Prefer server-reported generation time so RTF is independent of client + # poll_interval_s sleep/overshoot baked into output.latency. + generation_s = output.video_generation_time_ms / 1000.0 + if generation_s > 0: + output.video_rtf = generation_s / output.video_duration + elif output.latency > 0: + output.video_rtf = output.latency / output.video_duration + + async def async_request_openai_chat_omni_completions( request_func_input: RequestFuncInput, session: aiohttp.ClientSession, @@ -1036,6 +1353,7 @@ async def async_request_openai_chat_omni_completions( output.image_generation_time_ms = 0.0 output.image_pixels = 0 output.denoise_step_latency_ms = 0.0 + output.peak_memory_mb = 0.0 completion_tokens_seen = 0 try: async with session.post(url=api_url, json=payload, headers=headers) as response: @@ -1065,6 +1383,7 @@ async def async_request_openai_chat_omni_completions( timestamp = time.perf_counter() data = json.loads(chunk) _update_output_stage_metrics_from_payload(output, data) + _update_output_peak_memory_from_payload(output, data) usage = data.get("usage") completion_tokens = None if isinstance(usage, dict): @@ -1281,16 +1600,243 @@ async def async_request_openai_chat_omni_completions( return output +def _finalize_image_json_http_response( + output: MixRequestFuncOutput, + *, + start_time: float, + status: int, + data: Mapping[str, object] | None, + error_text: str | None, +) -> None: + """Set e2el after the image JSON body has been fully read and validated.""" + output.latency = time.perf_counter() - start_time + if status != 200: + output.error = f"HTTP {status}: {error_text or ''}" + output.success = False + return + if not isinstance(data, Mapping): + output.error = "HTTP 200 response did not contain a JSON object" + output.success = False + return + payload_image_count = _apply_image_metrics_from_payload(output, data) + if payload_image_count <= 0: + output.error = "HTTP 200 response did not contain a valid image payload" + output.success = False + return + output.success = True + + +async def async_request_openai_image_generations_omni( + request_func_input: RequestFuncInput, + session: aiohttp.ClientSession, + pbar: tqdm | None = None, +) -> MixRequestFuncOutput: + """JSON request to /v1/images/generations for image generation benchmarks.""" + api_url = request_func_input.api_url + _validate_api_url(api_url, "OpenAI Image Generations API", "images/generations") + + extra_body = dict(request_func_input.extra_body or {}) + model = request_func_input.model_name if request_func_input.model_name else request_func_input.model + output = MixRequestFuncOutput() + output.prompt_len = request_func_input.prompt_len + output.itl = [] + output.stage_metrics = {} + output.output_tokens = 0 + output.image_count = 0 + output.image_generation_time_ms = 0.0 + output.image_pixels = 0 + output.denoise_step_latency_ms = 0.0 + + size = extra_body.get("size") + if size is None: + width, height = extra_body.get("width"), extra_body.get("height") + if width is not None and height is not None: + size = f"{width}x{height}" + + payload: dict[str, object] = { + "model": model, + "prompt": request_func_input.prompt, + "n": int(extra_body.pop("n", extra_body.pop("num_outputs_per_prompt", 1)) or 1), + "response_format": "b64_json", + } + if size is not None: + payload["size"] = str(size) + + for key, value in extra_body.items(): + if key in {"height", "width"}: + continue + payload.setdefault(key, value) + + headers = { + "Content-Type": "application/json", + "Authorization": f"Bearer {os.environ.get('OPENAI_API_KEY')}", + } + _update_headers_common(headers, request_func_input) + + st = time.perf_counter() + output.start_time = st + try: + async with session.post(url=api_url, json=payload, headers=headers) as response: + if response.status == 200: + data = await response.json() + _finalize_image_json_http_response( + output, + start_time=st, + status=response.status, + data=data if isinstance(data, Mapping) else None, + error_text=None, + ) + else: + error_text = await response.text() + _finalize_image_json_http_response( + output, + start_time=st, + status=response.status, + data=None, + error_text=error_text, + ) + except Exception: + output.latency = time.perf_counter() - st + output.success = False + output.error = traceback.format_exc() + logger.error(f"ERROR: send image generation request failed, reason is: {output.error}") + + if pbar: + pbar.update(1) + return output + + +async def async_request_openai_videos_omni( + request_func_input: RequestFuncInput, + session: aiohttp.ClientSession, + pbar: tqdm | None = None, +) -> MixRequestFuncOutput: + """Multipart request to async /v1/videos, polling metadata until completion.""" + api_url = request_func_input.api_url + _validate_api_url(api_url, "OpenAI Videos API", "videos") + + extra_body = dict(request_func_input.extra_body or {}) + model = request_func_input.model_name if request_func_input.model_name else request_func_input.model + output = MixRequestFuncOutput() + output.prompt_len = request_func_input.prompt_len + output.itl = [] + output.stage_metrics = {} + output.output_tokens = 0 + output.video_duration = 0.0 + output.video_generation_time_ms = 0.0 + + request_body: dict[str, object] = { + "model": model, + "prompt": request_func_input.prompt, + } + request_body.update(extra_body) + + size = request_body.get("size") + if size is None: + width, height = request_body.get("width"), request_body.get("height") + if width is not None and height is not None: + size = f"{width}x{height}" + if size is not None: + request_body["size"] = str(size) + + form = aiohttp.FormData() + form.add_field("model", str(model)) + form.add_field("prompt", str(request_func_input.prompt)) + if request_body.get("size") is not None: + form.add_field("size", str(request_body["size"])) + _add_video_extra_body_to_form(form, extra_body, request_body) + + reference_added = False + for reference in _iter_image_edit_inputs(request_func_input.multi_modal_content): + if _add_video_reference_to_form(form, reference): + reference_added = True + break + if not reference_added: + image_reference = extra_body.get("image_reference") + if image_reference is not None: + _add_video_reference_to_form(form, image_reference) + + headers = { + "Authorization": f"Bearer {os.environ.get('OPENAI_API_KEY')}", + } + _update_headers_common(headers, request_func_input) + + poll_interval_s = float(extra_body.get("poll_interval_s", 2.0) or 2.0) + timeout_s = float(extra_body.get("poll_timeout_s", 6 * 60 * 60) or (6 * 60 * 60)) + st = time.perf_counter() + output.start_time = st + try: + async with session.post(url=api_url, data=form, headers=headers) as response: + if response.status != 200: + output.latency = time.perf_counter() - st + output.error = f"HTTP {response.status}: {await response.text()}" + output.success = False + return output + create_payload = await response.json() + + job_id = create_payload.get("id") + job_status = create_payload.get("status") + if not isinstance(job_id, str) or not job_id: + output.latency = time.perf_counter() - st + output.error = "Video creation response missing job id." + output.success = False + return output + + job_url = f"{api_url.rstrip('/')}/{job_id}" + poll_payload = create_payload + deadline = time.perf_counter() + timeout_s + while job_status not in {"completed", "failed"}: + if time.perf_counter() >= deadline: + output.latency = time.perf_counter() - st + output.error = f"Timed out waiting for video job {job_id} to complete." + output.success = False + return output + await asyncio.sleep(poll_interval_s) + async with session.get(job_url, headers=headers) as poll_response: + if poll_response.status != 200: + output.latency = time.perf_counter() - st + output.error = f"Polling failed HTTP {poll_response.status}: {await poll_response.text()}" + output.success = False + return output + poll_payload = await poll_response.json() + job_status = poll_payload.get("status") + + output.latency = time.perf_counter() - st + if job_status == "failed": + output.error = f"Video job failed: {poll_payload}" + output.success = False + return output + + _apply_video_metrics_from_payload(output, poll_payload, request_body) + output.success = True + except Exception: + output.latency = time.perf_counter() - st + output.success = False + output.error = traceback.format_exc() + logger.error(f"ERROR: send video request failed, reason is: {output.error}") + finally: + if pbar: + pbar.update(1) + return output + + async def async_request_openai_image_edits_omni( request_func_input: RequestFuncInput, session: aiohttp.ClientSession, pbar: tqdm | None = None, ) -> MixRequestFuncOutput: - """Streaming request to /v1/images/edits for multi-stage image-edit benchmarks.""" + """Multipart request to /v1/images/edits. + + Defaults to non-streaming JSON so single-stage edit models work. The server + rejects ``stream=true`` when ``len(stage_configs) <= 1``. Pass + ``stream: true`` in ``--extra-body`` for multi-stage SSE (AR TTFT / image + chunks). + """ api_url = request_func_input.api_url _validate_api_url(api_url, "OpenAI Image Edits API", "images/edits") extra_body = dict(request_func_input.extra_body or {}) + want_stream = coerce_bool(extra_body.pop("stream", None), default=False) model = request_func_input.model_name if request_func_input.model_name else request_func_input.model output = MixRequestFuncOutput() output.prompt_len = request_func_input.prompt_len @@ -1307,7 +1853,7 @@ async def async_request_openai_image_edits_omni( form.add_field("prompt", request_func_input.prompt) form.add_field("response_format", "b64_json") form.add_field("output_format", str(extra_body.get("output_format", "png"))) - form.add_field("stream", "true") + form.add_field("stream", "true" if want_stream else "false") size = extra_body.get("size") if size is None: @@ -1340,12 +1886,21 @@ async def async_request_openai_image_edits_omni( st = time.perf_counter() output.start_time = st - timestamp = st - most_recent_text_timestamp = st - generated_text = "" try: async with session.post(url=api_url, data=form, headers=headers) as response: - if response.status == 200: + if response.status != 200: + error_text = await response.text() + _finalize_image_json_http_response( + output, + start_time=st, + status=response.status, + data=None, + error_text=error_text, + ) + elif want_stream: + timestamp = st + most_recent_text_timestamp = st + generated_text = "" handler = StreamedResponseHandler() async for chunk_bytes in response.content.iter_any(): if not chunk_bytes: @@ -1366,6 +1921,7 @@ async def async_request_openai_image_edits_omni( data, update_output_tokens=(data.get("type") == "ar_delta"), ) + _update_output_peak_memory_from_payload(output, data) chunk_type = data.get("type") if chunk_type == "ar_delta": @@ -1401,9 +1957,16 @@ async def async_request_openai_image_edits_omni( output.generated_text = generated_text output.success = True else: - output.error = f"HTTP {response.status}: {await response.text()}" - output.success = False + data = await response.json() + _finalize_image_json_http_response( + output, + start_time=st, + status=response.status, + data=data if isinstance(data, Mapping) else None, + error_text=None, + ) except Exception: + output.latency = time.perf_counter() - st output.success = False output.error = traceback.format_exc() logger.error(f"ERROR: send image edit request failed, reason is: {output.error}") @@ -1951,6 +2514,18 @@ async def async_request_openai_realtime_duplex( if "openai-chat-omni" not in OPENAI_COMPATIBLE_BACKENDS: OPENAI_COMPATIBLE_BACKENDS.append("openai-chat-omni") +ASYNC_REQUEST_FUNCS["/v1/images/edits"] = async_request_openai_image_edits_omni +if "/v1/images/edits" not in OPENAI_COMPATIBLE_BACKENDS: + OPENAI_COMPATIBLE_BACKENDS.append("/v1/images/edits") + +ASYNC_REQUEST_FUNCS["/v1/images/generations"] = async_request_openai_image_generations_omni +if "/v1/images/generations" not in OPENAI_COMPATIBLE_BACKENDS: + OPENAI_COMPATIBLE_BACKENDS.append("/v1/images/generations") + +ASYNC_REQUEST_FUNCS["/v1/videos"] = async_request_openai_videos_omni +if "/v1/videos" not in OPENAI_COMPATIBLE_BACKENDS: + OPENAI_COMPATIBLE_BACKENDS.append("/v1/videos") + ASYNC_REQUEST_FUNCS["openai-audio-speech"] = async_request_openai_audio_speech if "openai-audio-speech" not in OPENAI_COMPATIBLE_BACKENDS: OPENAI_COMPATIBLE_BACKENDS.append("openai-audio-speech") @@ -2330,6 +2905,18 @@ def measured_ttft(output: RequestFuncOutput) -> float | None: defs.IMAGE_THROUGHPUT: getattr(metrics, defs.IMAGE_THROUGHPUT), defs.AVERAGE_PIXELS_PER_IMAGE: getattr(metrics, defs.AVERAGE_PIXELS_PER_IMAGE), defs.MEAN_DENOISE_STEP_LATENCY_MS: getattr(metrics, defs.MEAN_DENOISE_STEP_LATENCY_MS), + defs.TOTAL_VIDEO_DURATION_S: getattr(metrics, defs.TOTAL_VIDEO_DURATION_S), + defs.TOTAL_VIDEO_FRAMES: getattr(metrics, defs.TOTAL_VIDEO_FRAMES), + defs.VIDEO_THROUGHPUT: getattr(metrics, defs.VIDEO_THROUGHPUT), + defs.MEAN_VIDEO_RTF: getattr(metrics, defs.MEAN_VIDEO_RTF), + defs.MEDIAN_VIDEO_RTF: getattr(metrics, defs.MEDIAN_VIDEO_RTF), + defs.PERCENTILES_VIDEO_RTF: getattr(metrics, defs.PERCENTILES_VIDEO_RTF), + defs.MEAN_VIDEO_GENERATION_MS: getattr(metrics, defs.MEAN_VIDEO_GENERATION_MS), + defs.MEDIAN_VIDEO_GENERATION_MS: getattr(metrics, defs.MEDIAN_VIDEO_GENERATION_MS), + defs.PERCENTILES_VIDEO_GENERATION_MS: getattr(metrics, defs.PERCENTILES_VIDEO_GENERATION_MS), + defs.MEAN_PEAK_MEMORY_MB: getattr(metrics, defs.MEAN_PEAK_MEMORY_MB), + defs.MEDIAN_PEAK_MEMORY_MB: getattr(metrics, defs.MEDIAN_PEAK_MEMORY_MB), + defs.PERCENTILES_PEAK_MEMORY_MB: getattr(metrics, defs.PERCENTILES_PEAK_MEMORY_MB), "input_lens": [output.prompt_len for output in outputs], "start_times": [output.start_time for output in outputs], "output_lens": actual_output_lens, diff --git a/vllm_omni/benchmarks/serve.py b/vllm_omni/benchmarks/serve.py index 2c4f8d13043..96878a990a5 100644 --- a/vllm_omni/benchmarks/serve.py +++ b/vllm_omni/benchmarks/serve.py @@ -23,6 +23,37 @@ should_request_stage_metrics, ) +_ENDPOINT_BACKEND_KEYS = { + "/v1/images/generations", + "/v1/images/edits", + "/v1/videos", +} + + +def _normalize_endpoint(endpoint: str | None) -> str | None: + if endpoint is None: + return None + endpoint = str(endpoint).strip() + if not endpoint: + return None + if not endpoint.startswith("/"): + endpoint = f"/{endpoint}" + return endpoint + + +def _use_endpoint_backend_when_implicit(args: argparse.Namespace) -> None: + explicit_keys = getattr(args, "explicit_keys", frozenset()) + if "backend" in explicit_keys: + return + # Upstream vLLM defaults --backend to "openai". Treat that non-explicit + # default the same as an empty backend so endpoint-driven omni handlers work. + backend = getattr(args, "backend", None) + if backend not in (None, "", "openai"): + return + endpoint = _normalize_endpoint(getattr(args, "endpoint", None)) + if endpoint in _ENDPOINT_BACKEND_KEYS: + args.backend = endpoint + def main(args: argparse.Namespace) -> dict[str, Any]: if getattr(args, "seed_tts_wer_eval", False): @@ -31,6 +62,7 @@ def main(args: argparse.Namespace) -> dict[str, Any]: os.environ["SEED_TTS_WER_SAVE_ITEMS"] = "1" if getattr(args, "daily_omni_save_eval_items", False): os.environ["DAILY_OMNI_SAVE_EVAL_ITEMS"] = "1" + _use_endpoint_backend_when_implicit(args) set_print_stage(getattr(args, "print_stage", False)) args.extra_body = maybe_enable_stage_metrics( getattr(args, "extra_body", None), diff --git a/vllm_omni/diffusion/models/helios/pipeline_helios.py b/vllm_omni/diffusion/models/helios/pipeline_helios.py index 45c6d069328..820f9a1d98a 100644 --- a/vllm_omni/diffusion/models/helios/pipeline_helios.py +++ b/vllm_omni/diffusion/models/helios/pipeline_helios.py @@ -1358,6 +1358,8 @@ def _stage1_sample( """Single-stage denoising loop for one chunk.""" batch_size = latents.shape[0] do_true_cfg = guidance_scale > 1.0 and negative_prompt_embeds is not None + # Distilled (DMD) needs the chunk-start noise for renoise; mirror step_scheduler. + stage1_start_latents = latents with self.progress_bar(total=len(timesteps)) as pbar: for i, t in enumerate(timesteps): @@ -1420,7 +1422,20 @@ def _stage1_sample( cfg_normalize=False, ) - latents = self.scheduler_step_maybe_with_cfg(noise_pred, t, latents, do_true_cfg) + if self.is_distilled: + latents = self.scheduler.step( + noise_pred, + t, + latents, + return_dict=False, + cur_sampling_step=i, + dmd_noisy_tensor=stage1_start_latents, + dmd_sigmas=self.scheduler.sigmas, + dmd_timesteps=self.scheduler.timesteps, + all_timesteps=timesteps, + )[0] + else: + latents = self.scheduler_step_maybe_with_cfg(noise_pred, t, latents, do_true_cfg) pbar.update() diff --git a/vllm_omni/engine/stage_pool.py b/vllm_omni/engine/stage_pool.py index e7de5f51cee..0ab6c1da731 100644 --- a/vllm_omni/engine/stage_pool.py +++ b/vllm_omni/engine/stage_pool.py @@ -739,10 +739,16 @@ def build_stage_metrics( def _infer_output_unit_type(self, request_outputs: list[Any], *, token_count: int) -> str: final_output_type = getattr(self.stage_client, "final_output_type", None) - if self._has_image_output(request_outputs) or final_output_type in {"image", "images"}: + # Prefer declared modality over payload heuristics: video diffusion often + # stores frames in ``images`` (see serving_video / output_formatter). + if final_output_type in {"video", "videos"}: + return "video" + if final_output_type in {"image", "images"}: return "image" - if self._has_video_output(request_outputs) or final_output_type in {"video", "videos"}: + if self._has_video_output(request_outputs): return "video" + if self._has_image_output(request_outputs): + return "image" if self._has_audio_output(request_outputs) or final_output_type == "audio": return "audio" if self._has_trajectory_latent_output(request_outputs) or self._has_latent_output(request_outputs): @@ -778,6 +784,10 @@ def _count_output_units( total_videos = sum(self._count_videos(ro) for ro in request_outputs) if total_videos > 0: return total_videos + # Video payloads are commonly carried on ``images`` for diffusion. + total_images = sum(self._count_images(ro) for ro in request_outputs) + if total_images > 0: + return total_images if unit_type == "latent": total_latents = sum( self._count_value_units(getattr(ro, "trajectory_latents", None)) diff --git a/vllm_omni/entrypoints/cli/benchmark/cli_args.py b/vllm_omni/entrypoints/cli/benchmark/cli_args.py index 2cf66876931..68cbb9f860d 100644 --- a/vllm_omni/entrypoints/cli/benchmark/cli_args.py +++ b/vllm_omni/entrypoints/cli/benchmark/cli_args.py @@ -103,8 +103,8 @@ def add_diffusion_cli_args(parser: argparse.ArgumentParser) -> None: type=str, default="think", help=( - "Default bot_task form field for --backend openai-image-edits-omni " - "(/v1/images/edits). " + "Default bot_task form field for image edits " + "(--backend openai-image-edits-omni or --endpoint /v1/images/edits). " 'Use --extra-body \'{"bot_task":"..."}\' to override per run.' ), ) @@ -352,6 +352,12 @@ def preprocess_serve_args(args: argparse.Namespace) -> None: raise ValueError("OmniInteract requires --max-concurrency to be positive") extra_body = dict(getattr(args, "extra_body", None) or {}) bot_task = getattr(args, "bot_task", None) - if getattr(args, "backend", None) == "openai-image-edits-omni" and bot_task is not None: + backend = getattr(args, "backend", None) + endpoint = getattr(args, "endpoint", None) + # serve.py remaps implicit backend to the endpoint path for image edits; + # inject bot_task for both the named backend and /v1/images/edits. + if bot_task is not None and ( + backend in ("openai-image-edits-omni", "/v1/images/edits") or endpoint == "/v1/images/edits" + ): extra_body.setdefault("bot_task", bot_task) args.extra_body = extra_body diff --git a/vllm_omni/entrypoints/openai/api_server.py b/vllm_omni/entrypoints/openai/api_server.py old mode 100755 new mode 100644 index b305cb0cd38..d49d95175b5 --- a/vllm_omni/entrypoints/openai/api_server.py +++ b/vllm_omni/entrypoints/openai/api_server.py @@ -15,7 +15,7 @@ import time import uuid from argparse import Namespace -from collections.abc import AsyncIterator, Mapping +from collections.abc import AsyncIterator, Mapping, Sequence from contextlib import asynccontextmanager, suppress from http import HTTPStatus from numbers import Integral @@ -131,6 +131,7 @@ from vllm_omni.entrypoints.openai.protocol.videos import ( SecondStr, SizeStr, + VideoAction, VideoDeleteResponse, VideoError, VideoGenerationRequest, @@ -1946,12 +1947,26 @@ async def show_available_models(raw_request: Request) -> JSONResponse: # Image generation API endpoints +def _build_image_response_metrics( + *, + response_metrics: Any, + stage_durations: Any, + peak_memory_mb: Any, +) -> dict[str, Any]: + """Merge detailed stage metrics with the legacy image timing fields.""" + metrics = dict(response_metrics) if isinstance(response_metrics, dict) else {} + metrics["stage_durations"] = stage_durations or None + metrics["peak_memory_mb"] = float(peak_memory_mb) if peak_memory_mb else None + return metrics + + def _build_image_generation_response( *, images: list[Image.Image], request: ImageGenerationRequest, stage_durations: Any, peak_memory_mb: Any, + response_metrics: Any = None, ) -> ImageGenerationResponse | StreamingResponse: """Encode generated images and apply the requested response format.""" output_format = _choose_output_format(request.output_format or "png", None) @@ -1966,10 +1981,11 @@ def _build_image_generation_response( "created": int(time.time()), "data": image_data, "output_format": output_format, - "metrics": { - "stage_durations": stage_durations or None, - "peak_memory_mb": float(peak_memory_mb) if peak_memory_mb else None, - }, + "metrics": _build_image_response_metrics( + response_metrics=response_metrics, + stage_durations=stage_durations, + peak_memory_mb=peak_memory_mb, + ), } if request.size is not None: response_kwargs["size"] = request.size @@ -2066,6 +2082,8 @@ async def generate_images( extra_body["use_system_prompt"] = request.use_system_prompt if request.system_prompt is not None: extra_body["system_prompt"] = request.system_prompt + if request.return_stage_metrics is not None: + extra_body["return_stage_metrics"] = request.return_stage_metrics generation_result = await chat_handler.generate_diffusion_images( prompt=request.prompt, @@ -2079,12 +2097,13 @@ async def generate_images( status_code=generation_result.error.code if generation_result.error else 400, content=generation_result.model_dump(), ) - flat_images, stage_durations, peak_memory_mb, _ = generation_result + flat_images, stage_durations, peak_memory_mb, _, response_metrics = generation_result return _build_image_generation_response( images=flat_images, request=request, stage_durations=stage_durations, peak_memory_mb=peak_memory_mb, + response_metrics=response_metrics, ) # Build params - pass through user values directly @@ -2174,11 +2193,13 @@ async def generate_images( stage_durations = getattr(result, "stage_durations", None) peak_memory_mb = getattr(result, "peak_memory_mb", None) + response_metrics = getattr(result, "metrics", None) if request.return_stage_metrics else None return _build_image_generation_response( images=images, request=request, stage_durations=stage_durations, peak_memory_mb=peak_memory_mb, + response_metrics=response_metrics, ) except (EngineGenerateError, EngineDeadError) as exc: @@ -2510,7 +2531,7 @@ async def edit_images( status_code=generation_result.error.code if generation_result.error else 400, detail=generation_result.message, ) - images, stage_durations, peak_memory_mb, cot_output = generation_result + images, stage_durations, peak_memory_mb, cot_output, response_metrics = generation_result else: # Single-stage diffusion: use the direct path. result = await _generate_with_async_omni( @@ -2523,6 +2544,7 @@ async def edit_images( images = _extract_images_from_result(result) stage_durations = getattr(result, "stage_durations", None) peak_memory_mb = getattr(result, "peak_memory_mb", None) + response_metrics = getattr(result, "metrics", None) if return_stage_metrics else None logger.debug(f"Successfully generated {len(images)} image(s)") @@ -2543,10 +2565,11 @@ async def edit_images( output_format=output_format, size=size_str, cot_output=cot_output, - metrics={ - "stage_durations": stage_durations or None, - "peak_memory_mb": float(peak_memory_mb) if peak_memory_mb else None, - }, + metrics=_build_image_response_metrics( + response_metrics=response_metrics, + stage_durations=stage_durations, + peak_memory_mb=peak_memory_mb, + ), ) except (EngineGenerateError, EngineDeadError) as exc: @@ -3010,12 +3033,22 @@ def _reference_video_decode_spec( def video_response_from_request(model_name: str, req: VideoGenerationRequest) -> VideoResponse: + video_params = req.resolve_video_params() + duration_s = None + if req.seconds is not None: + duration_s = float(req.seconds) + elif video_params.num_frames is not None and video_params.fps is not None: + duration_s = video_params.num_frames / video_params.fps + resp = VideoResponse( model=model_name, status=VideoGenerationStatus.QUEUED, size=req.size, prompt=req.prompt, quality=req.quality or "default", + fps=video_params.fps, + num_frames=video_params.num_frames, + duration_s=duration_s, ) resp.seconds = str(req.seconds or resp.seconds) return resp @@ -3091,6 +3124,25 @@ def _cleanup_video_references( os.unlink(control_path) +def _unpack_video_generation_result( + result: Sequence[object], +) -> tuple[bytes, dict[str, float], float, VideoAction | None, dict[str, object]]: + video_metadata: dict[str, object] = {} + if len(result) == 5: + video_bytes, stage_durations, peak_memory_mb, action, raw_metadata = result + if isinstance(raw_metadata, dict): + video_metadata = {str(key): value for key, value in raw_metadata.items()} + else: + video_bytes, stage_durations, peak_memory_mb, action = result + return ( + cast(bytes, video_bytes), + cast(dict[str, float], stage_durations), + float(cast(float, peak_memory_mb)), + cast(VideoAction | None, action), + video_metadata, + ) + + async def _run_video_generation_job( handler: OmniOpenAIServingVideo, request: VideoGenerationRequest, @@ -3110,12 +3162,14 @@ async def _run_video_generation_job( await VIDEO_STORE.update_fields(video_id, {"status": VideoGenerationStatus.IN_PROGRESS}) started_at = time.perf_counter() try: - video_bytes, stage_durations, peak_memory_mb, action = await handler.generate_video_bytes( - request, - video_id, - reference_image=reference_image, - reference_video=reference_video, - reference_audio=reference_audio, + video_bytes, stage_durations, peak_memory_mb, action, video_metadata = _unpack_video_generation_result( + await handler.generate_video_bytes( + request, + video_id, + reference_image=reference_image, + reference_video=reference_video, + reference_audio=reference_audio, + ) ) save_context = await STORAGE_MANAGER.save(video_bytes, video_id) @@ -3131,6 +3185,7 @@ async def _run_video_generation_job( "peak_memory_mb": peak_memory_mb, "action": action, } + updated_fields.update({key: value for key, value in video_metadata.items() if value is not None}) if save_context.expires_at is not None: updated_fields["expires_at"] = save_context.expires_at @@ -3494,6 +3549,7 @@ async def _parse_video_form( frame_interpolation_model_path: str | None = Form(default=None), lora: str | None = Form(default=None), extra_params: str | None = Form(default=None), + return_stage_metrics: bool | None = Form(default=None), ) -> tuple[ VideoGenerationRequest, "OmniOpenAIServingVideo", @@ -3564,6 +3620,7 @@ async def _parse_video_form( "frame_interpolation_model_path": frame_interpolation_model_path, "lora": _parse_form_json(lora, expected_type=dict), "extra_params": _parse_form_json(extra_params, expected_type=dict), + "return_stage_metrics": return_stage_metrics, } request_data = {k: v for k, v in request_data.items() if v is not None} request = VideoGenerationRequest(**request_data) @@ -3844,15 +3901,17 @@ async def create_video_sync( raw_request.state.request_metadata = RequestResponseMetadata(request_id=request_id) started_at = time.perf_counter() try: - video_bytes, stage_durations, peak_memory_mb, _action = await asyncio.wait_for( - handler.generate_video_bytes( - request, - request_id, - reference_image=reference_image, - reference_video=reference_video, - reference_audio=reference_audio, + video_bytes, stage_durations, peak_memory_mb, _action, _video_metadata = _unpack_video_generation_result( + await asyncio.wait_for( + handler.generate_video_bytes( + request, + request_id, + reference_image=reference_image, + reference_video=reference_video, + reference_audio=reference_audio, + ), + timeout=VIDEO_SYNC_TIMEOUT_S, ), - timeout=VIDEO_SYNC_TIMEOUT_S, ) except asyncio.TimeoutError: raise HTTPException( diff --git a/vllm_omni/entrypoints/openai/protocol/images.py b/vllm_omni/entrypoints/openai/protocol/images.py index 2f28b6199e6..0215e1ce0b0 100644 --- a/vllm_omni/entrypoints/openai/protocol/images.py +++ b/vllm_omni/entrypoints/openai/protocol/images.py @@ -175,6 +175,10 @@ def validate_use_system_prompt(cls, v): default=None, description="Output image format: 'png', 'jpeg', or 'webp'. Defaults to 'png'.", ) + return_stage_metrics: bool | None = Field( + default=None, + description="Return stage metrics for benchmark clients.", + ) class ImageData(BaseModel): diff --git a/vllm_omni/entrypoints/openai/protocol/videos.py b/vllm_omni/entrypoints/openai/protocol/videos.py index fb4c395d709..cb8cc198d2c 100644 --- a/vllm_omni/entrypoints/openai/protocol/videos.py +++ b/vllm_omni/entrypoints/openai/protocol/videos.py @@ -290,6 +290,10 @@ def validate_quality(cls, value: str | None) -> str | None: default=None, description=("Optional model-specific parameters passed directly to the model's extra_args. "), ) + return_stage_metrics: bool | None = Field( + default=None, + description="Whether to include server-side stage metrics in async video metadata.", + ) def resolve_video_params( self, @@ -407,6 +411,13 @@ class VideoResponse(BaseModel): description="Filename of the saved output video files for this job.", ) inference_time_s: float | None = Field(default=None, description="End-to-end inference time in seconds.") + fps: float | None = Field(default=None, description="Resolved output video frames per second, if known.") + num_frames: int | None = Field(default=None, description="Resolved number of output video frames, if known.") + duration_s: float | None = Field(default=None, description="Resolved output video duration in seconds, if known.") + metrics: dict[str, Any] | None = Field( + default=None, + description="Optional profiler and stage metrics for benchmark clients.", + ) stage_durations: dict[str, float] = Field( default_factory=dict, description="Profiler stage durations reported by the diffusion pipeline.", diff --git a/vllm_omni/entrypoints/openai/serving_chat.py b/vllm_omni/entrypoints/openai/serving_chat.py old mode 100755 new mode 100644 index e799fce8141..bfa241a40ea --- a/vllm_omni/entrypoints/openai/serving_chat.py +++ b/vllm_omni/entrypoints/openai/serving_chat.py @@ -3288,7 +3288,11 @@ async def generate_diffusion_images( output_compression: int = 100, size: str = "auto", raw_request: Request | None = None, - ) -> tuple[list[Image.Image], dict[str, Any], float, str | None] | ErrorResponse | AsyncIterator[str]: + ) -> ( + tuple[list[Image.Image], dict[str, Any], float, str | None, dict[str, Any] | None] + | ErrorResponse + | AsyncIterator[str] + ): """Generate diffusion images and return raw images plus generation stats.""" if request_id is None: request_id = f"chatcmpl-{uuid.uuid4().hex[:16]}" @@ -3375,6 +3379,7 @@ async def generate_diffusion_images( images = getattr(result, "images", []) stage_durations = result.stage_durations peak_memory_mb = result.peak_memory_mb + response_metrics = getattr(result, "metrics", None) if return_stage_metrics else None cot_output = None req_out = result @@ -3407,7 +3412,7 @@ async def generate_diffusion_images( if isinstance(ar_text, str) and ar_text.strip(): cot_output = ar_text - return self._flatten_diffusion_images(images), stage_durations, peak_memory_mb, cot_output + return self._flatten_diffusion_images(images), stage_durations, peak_memory_mb, cot_output, response_metrics async def _stream_diffusion_image_chunks( self, diff --git a/vllm_omni/entrypoints/openai/serving_video.py b/vllm_omni/entrypoints/openai/serving_video.py index 939d55ef8b2..b3afe8b40b7 100644 --- a/vllm_omni/entrypoints/openai/serving_video.py +++ b/vllm_omni/entrypoints/openai/serving_video.py @@ -37,6 +37,7 @@ encode_video_base64, ) from vllm_omni.inputs.data import OmniDiffusionSamplingParams, OmniTextPrompt +from vllm_omni.metrics import count_video_frames from vllm_omni.model_extras import get_video_generation_defaults, should_preserve_reference_image_size from vllm_omni.model_extras.video_generation import VideoGenerationDefaults from vllm_omni.outputs.output_metadata import ( @@ -87,6 +88,27 @@ class VideoGenerationArtifacts: output_fps: float stage_durations: dict[str, float] peak_memory_mb: float + metrics: dict[str, object] | None = None + + +def _video_metadata_from_artifacts(artifacts: VideoGenerationArtifacts) -> dict[str, object]: + metadata: dict[str, object] = {} + if artifacts.output_fps > 0: + metadata["fps"] = artifacts.output_fps + + if not artifacts.videos: + return metadata + + num_frames = count_video_frames(artifacts.videos[0]) + if num_frames is not None and num_frames > 0: + metadata["num_frames"] = num_frames + if artifacts.output_fps > 0: + metadata["duration_s"] = num_frames / artifacts.output_fps + + if artifacts.metrics: + metadata["metrics"] = artifacts.metrics + + return metadata class OmniOpenAIServingVideo: @@ -422,6 +444,8 @@ async def _run_and_extract( model_fps = self._resolve_fps(result) output_fps_base = (vp.fps if fps_provided else None) or model_fps or vp.fps or 24 output_fps = output_fps_base * self._resolve_video_fps_multiplier(result) + raw_metrics = getattr(result, "metrics", None) if request.return_stage_metrics else None + metrics = {str(key): value for key, value in raw_metrics.items()} if isinstance(raw_metrics, Mapping) else None return VideoGenerationArtifacts( videos=videos, audios=audios, @@ -430,6 +454,7 @@ async def _run_and_extract( output_fps=output_fps, stage_durations=self._extract_stage_durations(result), peak_memory_mb=self._extract_peak_memory_mb(result), + metrics=metrics, ) async def generate_videos( @@ -495,7 +520,7 @@ async def generate_video_bytes( reference_image: ReferenceImage | None = None, reference_video: ReferenceVideo | None = None, reference_audio: ReferenceAudio | None = None, - ) -> tuple[bytes, dict[str, float], float, VideoAction | None]: + ) -> tuple[bytes, dict[str, float], float, VideoAction | None, dict[str, object]]: """Generate a video and return raw MP4 bytes, bypassing base64 encoding.""" artifacts = await self._run_and_extract( request, @@ -518,9 +543,10 @@ async def generate_video_bytes( video_codec_options = request.extra_params["video_codec_options"] action = artifacts.actions[0] + video_metadata = _video_metadata_from_artifacts(artifacts) if action is not None and isinstance(artifacts.videos[0], dict): logger.info("Action-only video request %s completed; skipping MP4 encoding.", reference_id) - return b"", artifacts.stage_durations, artifacts.peak_memory_mb, action + return b"", artifacts.stage_durations, artifacts.peak_memory_mb, action, video_metadata _t_encode_start = time.perf_counter() video_bytes = _encode_video_bytes( @@ -532,7 +558,7 @@ async def generate_video_bytes( ) _t_encode_ms = (time.perf_counter() - _t_encode_start) * 1000 logger.info("Video response encoding (MP4 bytes): %.2f ms", _t_encode_ms) - return video_bytes, artifacts.stage_durations, artifacts.peak_memory_mb, artifacts.actions[0] + return video_bytes, artifacts.stage_durations, artifacts.peak_memory_mb, artifacts.actions[0], video_metadata @staticmethod def _resolve_video_fps_multiplier(result: object) -> int: diff --git a/vllm_omni/metrics/__init__.py b/vllm_omni/metrics/__init__.py index 7b6a2979f55..c510c943b4c 100644 --- a/vllm_omni/metrics/__init__.py +++ b/vllm_omni/metrics/__init__.py @@ -1,3 +1,6 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project + from .prometheus import OmniPrometheusMetrics, OmniRequestCounter from .stats import OrchestratorAggregator, StageRequestStats, StageStats from .utils import ( @@ -5,6 +8,7 @@ count_audio_frames, count_image_pixels, count_tokens_from_outputs, + count_video_frames, ) __all__ = [ @@ -17,4 +21,5 @@ "count_audio_frames", "count_image_pixels", "count_tokens_from_outputs", + "count_video_frames", ] diff --git a/vllm_omni/metrics/definitions.py b/vllm_omni/metrics/definitions.py index 68adb44ca49..4b0043bad71 100644 --- a/vllm_omni/metrics/definitions.py +++ b/vllm_omni/metrics/definitions.py @@ -76,6 +76,29 @@ STD_IMAGE_GENERATION_MS = f"std_{IMAGE_GENERATION}_ms" PERCENTILES_IMAGE_GENERATION_MS = f"percentiles_{IMAGE_GENERATION}_ms" +VIDEO_DURATION = "video_duration" +VIDEO_RTF = "video_rtf" +VIDEO_GENERATION = "video_generation" +VIDEO_GENERATION_TIME_MS = f"{VIDEO_GENERATION}_time_ms" +VIDEO_FRAMES = "video_frames" +TOTAL_VIDEO_DURATION_S = f"total_{VIDEO_DURATION}_s" +TOTAL_VIDEO_FRAMES = f"total_{VIDEO_FRAMES}" +VIDEO_THROUGHPUT = "video_throughput" +MEAN_VIDEO_RTF = f"mean_{VIDEO_RTF}" +MEDIAN_VIDEO_RTF = f"median_{VIDEO_RTF}" +STD_VIDEO_RTF = f"std_{VIDEO_RTF}" +PERCENTILES_VIDEO_RTF = f"percentiles_{VIDEO_RTF}" +MEAN_VIDEO_GENERATION_MS = f"mean_{VIDEO_GENERATION}_ms" +MEDIAN_VIDEO_GENERATION_MS = f"median_{VIDEO_GENERATION}_ms" +STD_VIDEO_GENERATION_MS = f"std_{VIDEO_GENERATION}_ms" +PERCENTILES_VIDEO_GENERATION_MS = f"percentiles_{VIDEO_GENERATION}_ms" + +PEAK_MEMORY_MB = "peak_memory_mb" +MEAN_PEAK_MEMORY_MB = f"mean_{PEAK_MEMORY_MB}" +MEDIAN_PEAK_MEMORY_MB = f"median_{PEAK_MEMORY_MB}" +STD_PEAK_MEMORY_MB = f"std_{PEAK_MEMORY_MB}" +PERCENTILES_PEAK_MEMORY_MB = f"percentiles_{PEAK_MEMORY_MB}" + # Stage snapshot / StageBenchmarkMetrics field names. TOTAL_OUTPUT = "total_output" TTFTS = "ttfts" @@ -166,7 +189,7 @@ NUM_INFERENCE_STEPS = METRIC_PREFIX + "num_inference_steps" IMAGE_COUNT_METRIC = METRIC_PREFIX + IMAGE_COUNT IMAGE_PIXELS_METRIC = METRIC_PREFIX + IMAGE_PIXELS -PEAK_MEMORY_MB = METRIC_PREFIX + "peak_memory_mb" +PEAK_MEMORY_MB_METRIC = METRIC_PREFIX + "peak_memory_mb" REQUESTS_FAILED = METRIC_PREFIX + "requests_failed" KV_WAIT_S = METRIC_PREFIX + "kv_wait_s" DIFFUSION_FORWARD_S = METRIC_PREFIX + "diffusion_forward_s" diff --git a/vllm_omni/metrics/prometheus.py b/vllm_omni/metrics/prometheus.py index 911d21e646a..2e68045c1ca 100644 --- a/vllm_omni/metrics/prometheus.py +++ b/vllm_omni/metrics/prometheus.py @@ -1,3 +1,6 @@ +# SPDX-License-Identifier: Apache-2.0 +# SPDX-FileCopyrightText: Copyright contributors to the vLLM-Omni project + from prometheus_client import Counter, Gauge, Histogram from vllm_omni.metrics import definitions as defs @@ -72,7 +75,7 @@ labelnames=_labelnames, ) _peak_memory_family = Gauge( - defs.PEAK_MEMORY_MB, + defs.PEAK_MEMORY_MB_METRIC, "Peak GPU memory in MB observed during a stage's generation.", labelnames=list(defs.DIFFUSION_LABELS), ) diff --git a/vllm_omni/metrics/utils.py b/vllm_omni/metrics/utils.py index 10cdab94a6c..ad58e4cec94 100644 --- a/vllm_omni/metrics/utils.py +++ b/vllm_omni/metrics/utils.py @@ -163,6 +163,24 @@ def extract_diffusion_denoise_ms(output: Any) -> float | None: return sum_diffusion_stage_durations_ms(output, ".diffuse") +def coerce_bool(value: object, *, default: bool = False) -> bool: + """Coerce JSON / form / env-like values to bool. + + Accepts ``bool``, numeric 0/1, and common truthy strings + (``\"1\"`` / ``\"true\"`` / ``\"yes\"`` / ``\"on\"``, case-insensitive). + ``None`` and unrecognized types return ``default``. + """ + if value is None: + return default + if isinstance(value, bool): + return value + if isinstance(value, (int, float)): + return bool(value) + if isinstance(value, str): + return value.strip().lower() in {"1", "true", "yes", "on"} + return default + + def coerce_positive_int_scalar(value: object) -> int | None: """Coerce a value to a positive int without importing tensor libs. @@ -199,6 +217,34 @@ def coerce_positive_int_scalar(value: object) -> int | None: return value if value > 0 else None +def coerce_positive_float_scalar(value: object) -> float | None: + """Coerce a scalar-like value to a positive float. + + Handles plain numeric values, one-element tensor/numpy-like values via + ``.item()``, and list/tuple wrappers. Returns ``None`` when no positive + float can be extracted. + """ + if value is None: + return None + if isinstance(value, (list, tuple)): + for item in value: + coerced = coerce_positive_float_scalar(item) + if coerced is not None: + return coerced + return None + item = getattr(value, "item", None) + if callable(item): + try: + value = item() + except Exception: + return None + try: + value = float(value) + except (TypeError, ValueError): + return None + return value if value > 0 else None + + def resolve_int_by_sequential_keys( source: Mapping[str, object] | object | None, keys: Sequence[str], @@ -370,6 +416,43 @@ def count_image_pixels(value: object) -> int: return dims[-2] * dims[-1] +def count_video_frames(video: object) -> int | None: + """Return the frame count for common nested and tensor video layouts.""" + if isinstance(video, Mapping): + for key in ("video", "frames", "images"): + if key in video: + return count_video_frames(video[key]) + return None + + if isinstance(video, (list, tuple)): + return len(video) + + shape = getattr(video, "shape", None) + if shape is None: + return None + try: + dims = tuple(int(dim) for dim in shape) + except (TypeError, ValueError): + return None + + if len(dims) == 5: + # Common layouts: [B, C, F, H, W] and [B, F, H, W, C]. + if dims[-1] in (1, 3, 4): + return dims[1] + if dims[1] in (1, 3, 4): + return dims[2] + return dims[1] + if len(dims) == 4: + # Common layouts: [F, H, W, C] and [C, F, H, W]. + if dims[-1] in (1, 3, 4): + return dims[0] + if dims[0] in (1, 3, 4): + return dims[1] + return dims[0] + + return None + + def _build_field_defs( cls: type, exclude: set[str],