Skip to content
Closed
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
107 changes: 107 additions & 0 deletions tests/engine/duplex/test_session_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
)
from vllm_omni.engine.duplex.session.manager import DuplexSessionManager
from vllm_omni.engine.duplex.session.runner import DuplexSessionRunner
from vllm_omni.engine.duplex.vad import SileroVADBackendProvider, SpeechDetectorBackend
from vllm_omni.metrics.stats import StageRequestStats, StageStats
from vllm_omni.model_executor.models.minicpmo_4_5.duplex.plugin import MiniCPMO45DuplexPlugin

Expand Down Expand Up @@ -199,6 +200,7 @@ async def open_harness(
runtime_config: DuplexSessionRuntimeConfig | None = None,
stage_count: int = 2,
clock: Any = None,
vad_backend_provider: SileroVADBackendProvider | None = None,
) -> Harness:
plugin = MiniCPMO45DuplexPlugin(_fake_encode_audio)
port = RecordingStagePort(stage_count=stage_count)
Expand All @@ -213,6 +215,8 @@ async def open_harness(
model_config=None,
clock=clock,
)
if vad_backend_provider is not None:
manager.vad_backend_provider = vad_backend_provider
body: dict[str, object] = {"auto_response": auto_response, **(extra_body or {})}
config = DuplexSessionConfig(
model="openbmb/MiniCPM-o-4_5",
Expand Down Expand Up @@ -1193,3 +1197,106 @@ async def test_server_vad_speech_stopped_still_commits_a_turn_mode_session() ->
assert len(_final_submissions(h)) == 1, "the detector's stop commits the turn and starts the response"
finally:
await close_harness(h)


class _ScriptedSileroBackend:
"""Scores each 512-sample frame from a script instead of running Silero."""

def __init__(self, probabilities: Sequence[float]) -> None:
self._probabilities = iter(probabilities)

def new_state(self) -> object:
return None

def infer(self, frame: np.ndarray, state: object) -> tuple[float, object]:
assert frame.size == 512
return next(self._probabilities, 0.0), state


class _ScriptedSileroProvider(SileroVADBackendProvider):
def __init__(self, probabilities: Sequence[float]) -> None:
super().__init__()
self._scripted = _ScriptedSileroBackend(probabilities)

def get(self) -> SpeechDetectorBackend:
return self._scripted


_SERVER_VAD = {
"type": "server_vad",
"threshold": 0.5,
"prefix_padding_ms": 300,
"silence_duration_ms": 500,
"min_speech_duration_ms": 96,
}


def _utterance_with_pause(pause_frames: int) -> list[float]:
"""Speech, a 32 ms dip, 640 ms of speech, a pause, 800 ms of speech, then 800 ms of silence.

The dip matters: it opens the silence candidate that the speech after it has
to cancel. If it is not cancelled, the pause that follows fires it at once.
"""
return [0.9] * 12 + [0.1] + [0.9] * 20 + [0.1] * pause_frames + [0.9] * 25 + [0.02] * 25


async def _append_as_200ms_chunks(h: Harness, frames: int) -> list[DuplexEvent]:
events: list[DuplexEvent] = []
remaining = frames * 512
while remaining:
size = min(3200, remaining)
remaining -= size
events += await h.run(append_audio(size, is_speech=None))
return events


@pytest.mark.asyncio
@pytest.mark.parametrize("pause_frames", [6, 11], ids=["192ms-pause", "352ms-pause"])
async def test_a_pause_shorter_than_the_silence_duration_does_not_commit_a_turn_mode_session(
pause_frames: int,
) -> None:
"""The real endpoint rules through the runner, not a scripted result.

With ``auto_response`` off the detector's stop is the commit, so an endpoint
fired inside the pause starts a response on the first half of the utterance.
"""
probabilities = _utterance_with_pause(pause_frames)
h = await open_harness(
auto_response=False,
extra_body={"realtime_turn_detection": dict(_SERVER_VAD)},
vad_backend_provider=_ScriptedSileroProvider(probabilities),
)
try:
events = await _append_as_200ms_chunks(h, len(probabilities))

speech_frames = len(probabilities) - 25
stops = [event for event in events if event.type == "input_audio_buffer.speech_stopped"]
assert types(events).count("input_audio_buffer.speech_started") == 1
# The end of the audio sent, 500 ms of trailing silence included.
assert [stop.audio_end_ms for stop in stops] == [(speech_frames + 16) * 32]
assert types(events).count("input_audio_buffer.committed") == 1
assert len(_final_submissions(h)) == 1
finally:
await close_harness(h)


@pytest.mark.asyncio
async def test_a_pause_shorter_than_the_silence_duration_keeps_an_auto_response_utterance_whole() -> None:
"""Nothing is committed either way here; an early endpoint splits the utterance in two instead.

The second segment opens with its own ``speech_started``, and that edge is
what ``barge_in_on_speech`` cancels a response on.
"""
probabilities = _utterance_with_pause(11)
h = await open_harness(
extra_body={"realtime_turn_detection": dict(_SERVER_VAD)},
vad_backend_provider=_ScriptedSileroProvider(probabilities),
)
try:
events = await _append_as_200ms_chunks(h, len(probabilities))

assert types(events).count("input_audio_buffer.speech_started") == 1
assert types(events).count("input_audio_buffer.speech_stopped") == 1
assert "input_audio_buffer.committed" not in types(events)
finally:
await close_harness(h)
149 changes: 141 additions & 8 deletions tests/engine/duplex/test_vad_backend.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,22 @@
"""The engine-side Silero VAD: backend selection and endpoint rules.

The endpoint rules are the part worth pinning hardest. They were reimplemented
when turn detection moved into the engine, and drifted from upstream's
``ThresholdEndpointPolicy`` in two ways that only show up mid-utterance: a loud
frame used to cancel an in-progress silence timer, and ``audio_end_ms`` used to
exclude the trailing silence OpenAI says it should include. The parity cases
below hold upstream's own outputs, so neither can come back unnoticed.
when turn detection moved into the engine, and the reimplementation has to keep
three bands apart, not two: below ``threshold - 0.15`` a turn may end, at or
above ``threshold`` the speaker is plainly still going and any pending endpoint
is cancelled, and in between the endpoint clock keeps running but cannot fire.
Collapse the top two bands and a breath pause followed by a second of clear
speech commits the turn on the next quiet frame. ``audio_end_ms`` is the other
rule that has drifted: OpenAI counts the trailing silence, an earlier version
did not. The parity cases below hold the outputs of the serving-side policy this
replaced, so neither can come back unnoticed.
"""

from __future__ import annotations

import base64
import hashlib
import itertools
import sys
import threading
from concurrent.futures import ThreadPoolExecutor
Expand Down Expand Up @@ -59,8 +64,8 @@ def _drive(config: SileroVADConfig, probabilities: list[float]) -> list[tuple[st

#: Probability sequences with the endpoint events upstream's
#: ``ThresholdEndpointPolicy`` produced for them, captured from
#: ``entrypoints/duplex/server_vad.py`` at `e2d2617f` before that module was
#: removed. Golden rather than computed because the oracle no longer ships: the
#: ``entrypoints/duplex/server_vad.py`` at ``99ff4f307~1``, the commit before
#: that module was removed. Golden rather than computed because the oracle no longer ships: the
#: serving layer is transport only, so the VAD it used to carry is gone. Any
#: drift in the engine's rules still breaks these.
_PARITY_CASES = {
Expand All @@ -74,11 +79,16 @@ def _drive(config: SileroVADConfig, probabilities: list[float]) -> list[tuple[st
{"threshold": 0.5, "prefix_padding_ms": 300, "silence_duration_ms": 64, "min_speech_duration_ms": 32},
[("start", 0), ("stop", 160)],
),
"loud frames inside the silence timer": (
"hysteresis-band frames inside the silence timer": (
[0.9, 0.9, 0.1, 0.45, 0.1, 0.1, 0.1],
{"threshold": 0.5, "prefix_padding_ms": 0, "silence_duration_ms": 96, "min_speech_duration_ms": 32},
[("start", 0), ("stop", 160)],
),
"clear speech inside the silence timer": (
[0.9, 0.9, 0.1, 0.9, 0.1, 0.1, 0.1, 0.1],
{"threshold": 0.5, "prefix_padding_ms": 0, "silence_duration_ms": 96, "min_speech_duration_ms": 32},
[("start", 0), ("stop", 224)],
),
"min speech duration rejects a blip": (
[0.9, 0.0, 0.9, 0.9, 0.9, 0.0, 0.0, 0.0, 0.0],
{"threshold": 0.5, "prefix_padding_ms": 0, "silence_duration_ms": 64, "min_speech_duration_ms": 96},
Expand Down Expand Up @@ -118,6 +128,44 @@ def test_hysteresis_band_does_not_restart_the_silence_timer() -> None:
assert _drive(config, [0.9, 0.9, 0.1, 0.45, 0.1, 0.1]) == [("start", 0), ("stop", 160)]


def test_the_reset_delays_an_endpoint_rather_than_removing_one() -> None:
"""Trailing silence still commits the turn, measured from the last clear frame.

A pause at frame 2 that is cancelled at frame 3 endpoints in the same place
as a turn that never paused: both have their last clear frame at index 3.
"""
config = SileroVADConfig(threshold=0.5, prefix_padding_ms=0, silence_duration_ms=96, min_speech_duration_ms=32)

paused = _drive(config, [0.9, 0.9, 0.1, 0.9, 0.1, 0.1, 0.1, 0.1])
uninterrupted = _drive(config, [0.9, 0.9, 0.9, 0.9, 0.1, 0.1, 0.1, 0.1])

assert paused == uninterrupted == [("start", 0), ("stop", 224)]


def test_a_reset_clears_a_candidate_that_clear_speech_had_already_cancelled() -> None:
"""Barge-in mid-turn: the silence clock goes with the rest of the stream state.

``reset()`` keeps the session clock, so the next turn's ``speech_start_ms``
still refers to the session timeline, but it must not inherit a silence
candidate from the turn it replaced.
"""
config = SileroVADConfig(threshold=0.5, prefix_padding_ms=0, silence_duration_ms=96, min_speech_duration_ms=32)
scores = iter([0.9, 0.1, 0.9, 0.9, 0.1, 0.1, 0.1])
vad = SileroStreamingVAD(config, frame_scorer=lambda _frame: next(scores))

for _ in range(4): # speech, a pause, then clear speech cancelling it
vad.process(np.zeros(FRAME, dtype=np.float32))
assert vad.speech_active
vad.reset()
assert not vad.speech_active

events = [vad.process(np.zeros(FRAME, dtype=np.float32)) for _ in range(3)]

# Nothing survives the reset: three quiet frames cannot end a turn that is
# no longer running, and no stale candidate fires one.
assert not any(result.speech_stopped for result in events)


def test_audio_end_ms_includes_the_trailing_silence() -> None:
"""OpenAI defines audio_end_ms as the end of the audio sent to the model.

Expand Down Expand Up @@ -213,6 +261,91 @@ def count(_frame: np.ndarray) -> float:
assert whole == split == 16_000 // FRAME


def _drive_in_chunks(
config: SileroVADConfig,
probabilities: list[float],
chunk_sizes: list[int],
*,
reset_at: int | None = None,
sample_rate_hz: int = SAMPLE_RATE_HZ,
) -> list[tuple[str, int | None]]:
"""``_drive``, but with the audio cut into ``chunk_sizes`` (cycled) the way appends arrive.

``reset_at`` is an input sample offset; the chunk spanning it is split there so
every cutting resets at the same point. Edges reported by one call are listed
in the order ``realtime_events`` emits them.
"""
scores = iter(probabilities)
vad = SileroStreamingVAD(config, frame_scorer=lambda _frame: next(scores, 0.0))
total = len(probabilities) * FRAME * sample_rate_hz // SAMPLE_RATE_HZ
events: list[tuple[str, int | None]] = []

def feed(size: int) -> None:
if size == 0:
return
silence = np.zeros(size, dtype=np.float32)
if sample_rate_hz == SAMPLE_RATE_HZ:
result = vad.process(silence)
else:
result = vad.process_base64(_pcm16_b64(silence), fmt="pcm16", sample_rate_hz=sample_rate_hz)
if result.speech_stopped and result.speech_active:
events.append(("stop", result.speech_end_ms))
if result.speech_started:
events.append(("start", result.speech_start_ms))
if result.speech_stopped and not result.speech_active:
events.append(("stop", result.speech_end_ms))

fed = 0
for size in itertools.cycle(chunk_sizes):
if fed == total:
return events
size = min(size, total - fed)
if reset_at is not None and fed < reset_at <= fed + size:
feed(reset_at - fed)
vad.reset()
feed(fed + size - reset_at)
reset_at = None
else:
feed(size)
fed += size
raise AssertionError("unreachable")


#: Speech, a 32 ms dip, 640 ms of speech, another dip, 800 ms of speech, then
#: 800 ms of silence. One segment, and only because clear speech cancels both dips.
_TWO_DIPS = [0.9] * 12 + [0.1] + [0.9] * 20 + [0.1] + [0.9] * 25 + [0.02] * 25


@pytest.mark.parametrize("prefix_padding_ms", [0, 300])
@pytest.mark.parametrize("reset", [False, True], ids=["no-reset", "mid-frame-reset"])
def test_endpoints_do_not_depend_on_where_the_audio_is_cut(prefix_padding_ms: int, reset: bool) -> None:
"""The engine scores whole appends; the endpoints must not move with the append size.

This is what ``tests/entrypoints/openai_api/test_server_vad.py`` checked for the
serving-side pipeline before #7413 removed it. One call reports at most one
start and one stop, so no cutting here is longer than a 200 ms append.
"""
config = SileroVADConfig(
threshold=0.5, prefix_padding_ms=prefix_padding_ms, silence_duration_ms=500, min_speech_duration_ms=96
)
reset_at = 20 * FRAME + 100 if reset else None
frame_by_frame = _drive_in_chunks(config, _TWO_DIPS, [FRAME], reset_at=reset_at)

for chunk_sizes in ([3_200], [73, 358, 1_023, 511, 513, 2_047], [7, 5, 11]):
assert _drive_in_chunks(config, _TWO_DIPS, chunk_sizes, reset_at=reset_at) == frame_by_frame, chunk_sizes
reset_at_24k = None if reset_at is None else reset_at * 3 // 2
at_24k = _drive_in_chunks(config, _TWO_DIPS, [110, 1_535, 4_799], reset_at=reset_at_24k, sample_rate_hz=24_000)
assert at_24k == frame_by_frame


def test_the_cut_up_utterance_ends_once_after_the_real_silence() -> None:
"""The case above does exercise the reset: without it the second dip ends the turn at 1088 ms."""
config = SileroVADConfig(threshold=0.5, prefix_padding_ms=300, silence_duration_ms=500, min_speech_duration_ms=96)

# 59 frames of speech, then 16 of silence to reach 500 ms: 75 * 32 ms = 2400 ms.
assert _drive_in_chunks(config, _TWO_DIPS, [3_200]) == [("start", 0), ("stop", 2_400)]


def test_a_sample_rate_change_mid_stream_is_rejected() -> None:
vad = SileroStreamingVAD(SileroVADConfig(), frame_scorer=lambda _frame: 0.0)
chunk = _pcm16_b64(np.zeros(1_024, dtype=np.float32))
Expand Down
Loading