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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 51 additions & 3 deletions tests/diffusion/test_diffusion_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -869,9 +869,45 @@ async def _consume_final_output(generator):
return final_output


@pytest.mark.cpu
@pytest.mark.asyncio
@pytest.mark.parametrize("output_wait", [0.0, 4.0])
async def test_step_streaming_excludes_output_wait_from_execution_time(
output_wait: float, mocker: MockerFixture
) -> None:
clock = [0.0]
mocker.patch.object(diffusion_engine_module.time, "perf_counter", side_effect=lambda: clock[0])
output = DiffusionOutput(stage_durations={"denoise": 10.0}, output_ready_wait_time=output_wait)
formatted_output = SimpleNamespace(metrics={})
request = SimpleNamespace(scheduler_queue_wait_ms=None)

async def output_stream(request_id):
# Ten seconds of execution followed by output materialization.
clock[0] = 10.0 + output_wait
yield output

engine = mocker.Mock(
pre_process_func=None,
_scheduler_num_waiting_reqs=0,
_check_and_start_background_loop=mocker.AsyncMock(),
_prepare_request_for_admission=mocker.Mock(return_value=request),
_add_prepared_request=mocker.Mock(return_value="timed-request"),
get_output_stream=output_stream,

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

[P2] test_step_streaming_excludes_output_wait_from_execution_time patches get_output…

Evidence and suggested fix

test_step_streaming_excludes_output_wait_from_execution_time patches get_output_stream and yields a DiffusionOutput already populated with output_ready_wait_time (tests/diffusion/test_diffusion_engine.py:880,895), so it only proves step_streaming subtraction. Production writes the field in get_output_stream after wait_output_ready (vllm_omni/diffusion/diffusion_engine.py:1177). The only real-stream check is test_async_output_is_claimed_after_materialization's >= 0.0 (tests/diffusion/test_diffusion_engine_cleanup.py:336), which also holds for the dataclass default 0.0 (vllm_omni/diffusion/data.py:1816). Deleting the assignment would not fail tests, and diffusion_engine_exec_time_ms would include the wait. Drive a non-zero wait through real get_output_stream (completed executor future + patched perf_counter) and assert output.output_ready_wait_time / output_ready_wait_time_ms and diffusion_engine_exec_time_ms.

Evidence: Trigger: the new timing test never exercises the production writer. tests/diffusion/test_diffusion_engine.py:880 output = DiffusionOutput(stage_durations={"denoise": 10.0}, output_ready_wait_time=output_wait) and :895 get_output_stream=output_stream, inject the field and mock the stream. Production assignment is vllm_omni/diffusion/diffusion_engine.py:1177 output.output_ready_wait_time = time.perf_counter() - output_ready_wait_start_time. Unmet requirement: tests/diffusion/test_diffusion_engine_cleanup.py:336 assert materialized_output.output_ready_wait_time >= 0.0 still passes if that assignment is deleted, because vllm_omni/diffusion/data.py:1816 output_ready_wait_time: float = 0.0 is the default — then step_streaming exec_total_time = time.perf_counter() - exec_start_time - output_ready_wait_time would count materialization wait inside diffusion_engine_exec_time_ms with no test failure.

postprocess_output=mocker.Mock(return_value=[formatted_output]),
)

results = [batch async for batch in DiffusionEngine.step_streaming(engine, request)]

assert results == [[formatted_output]]
assert formatted_output.metrics["diffusion_engine_exec_time_ms"] == pytest.approx(10_000.0)
assert formatted_output.metrics["output_ready_wait_time_ms"] == pytest.approx(output_wait * 1000)
assert formatted_output.metrics["postprocess_time_ms"] == 0.0


@pytest.mark.cpu
@pytest.mark.parametrize("entrypoint", ["add_request", "async_add_req_and_stream_response"])
def test_engine_admission_preprocesses_request_once(entrypoint: str) -> None:
@pytest.mark.asyncio
async def test_engine_admission_preprocesses_request_once(entrypoint: str, mocker: MockerFixture) -> None:
raw_request = OmniDiffusionRequest(
prompt="raw",
sampling_params=OmniDiffusionSamplingParams(num_inference_steps=1),
Expand All @@ -891,7 +927,17 @@ def preprocess(request):

engine = _make_admission_engine(preprocess)

getattr(engine, entrypoint)(raw_request)
if entrypoint == "async_add_req_and_stream_response":
output = DiffusionOutput(output="prepared", finished=True)

async def output_stream(request_id):
assert request_id == prepared_request.request_id
yield output

mocker.patch.object(engine, "get_output_stream", side_effect=output_stream)
assert [result async for result in engine.async_add_req_and_stream_response(raw_request)] == [output]
else:
engine.add_request(raw_request)

assert preprocess_calls == [raw_request]
assert engine.scheduler._waiting_queue == [prepared_request]
Expand All @@ -903,6 +949,8 @@ async def test_async_add_req_and_stream_response():
engine = object.__new__(DiffusionEngine)
engine.scheduler = MockScheduler()
engine._out_streams = {}
engine._unclaimed_async_outputs = {}
engine._shutdown_output_futures = {}
engine.abort_queue: queue.Queue[str] = queue.Queue()
engine._rpc_queue = queue.Queue()
engine._rpc_lock = threading.RLock()
Expand All @@ -924,7 +972,7 @@ async def test_async_add_req_and_stream_response():

def _finalize(rid, out, err=None, **kwargs):
# Stream consumers stop on ``finished``; keep result_data for assertions.
return SimpleNamespace(result_data=out.result.result_data, finished=True)
return SimpleNamespace(result_data=out.result.result_data, finished=True, async_output_id=None)

engine._finalize_finished_request = _finalize

Expand Down
Loading
Loading