3277 lines
130 KiB
Python
3277 lines
130 KiB
Python
import asyncio
|
|
import contextlib
|
|
import struct
|
|
import threading
|
|
import time
|
|
import uuid
|
|
|
|
from services.ai_voice_runtime_service.audiosocket import (
|
|
AUDIO_SOCKET_PACKET_PCM16,
|
|
AUDIO_SOCKET_PACKET_UUID,
|
|
EnergyVAD,
|
|
encode_audio_packet,
|
|
encode_packet,
|
|
normalize_media_uuid,
|
|
read_packet,
|
|
resample_pcm16le,
|
|
)
|
|
from services.ai_voice_runtime_service.media_runtime import AudioSocketMediaRuntime, MediaActor, MediaRegistration
|
|
from services.ai_voice_runtime_service.providers.asr import (
|
|
ASRProvider,
|
|
ASRTranscription,
|
|
StreamingASRPartial,
|
|
StreamingASRProvider,
|
|
StreamingASRUnavailable,
|
|
)
|
|
from services.ai_voice_runtime_service.providers.tts import TTSProvider, TTSSynthesis
|
|
from services.shared.models import VoiceAITurnDecisionOut
|
|
|
|
|
|
class _StubASRProvider(ASRProvider):
|
|
name = "stub-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="hello", language=language_hint or "ru", confidence=0.9)
|
|
|
|
|
|
class _StubTTSProvider(TTSProvider):
|
|
name = "stub-tts"
|
|
|
|
def synthesize(
|
|
self,
|
|
text: str,
|
|
*,
|
|
language: str | None = None,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> TTSSynthesis:
|
|
del language, style_hints
|
|
assert text
|
|
return TTSSynthesis(
|
|
text=text,
|
|
audio_bytes=(b"\x10\x00" * 960),
|
|
sample_rate_hz=24000,
|
|
)
|
|
|
|
|
|
def test_resample_pcm16le_downsamples_to_8khz():
|
|
source = b"\x20\x00" * 2400
|
|
converted = resample_pcm16le(source, input_rate_hz=24000, output_rate_hz=8000)
|
|
assert converted
|
|
assert len(converted) < len(source)
|
|
|
|
|
|
def test_normalize_media_uuid_accepts_text_bytes():
|
|
media_uuid = str(uuid.uuid4())
|
|
assert normalize_media_uuid(media_uuid.encode("utf-8")) == media_uuid
|
|
|
|
|
|
def test_energy_vad_emits_utterance_after_trailing_silence():
|
|
vad = EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400)
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
|
|
first = vad.feed(speech_frame)
|
|
second = vad.feed(speech_frame)
|
|
third = vad.feed(silence_frame)
|
|
fourth = vad.feed(silence_frame)
|
|
|
|
assert first.speech_started is False
|
|
assert second.speech_started is True
|
|
assert third.utterance_pcm is None
|
|
assert fourth.utterance_pcm is not None
|
|
|
|
|
|
def test_media_runtime_streams_greeting_and_turn():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
states: list[tuple[str, str]] = []
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
turns: list[tuple[str, str, bool]] = []
|
|
touch_calls: list[str] = []
|
|
handoffs: list[str] = []
|
|
media_ended: list[tuple[str, str]] = []
|
|
errors: list[tuple[str, str]] = []
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime",
|
|
call_id="call_media_runtime",
|
|
interaction_id="int_media_runtime",
|
|
ai_session_id="ais_media_runtime",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: media_ended.append((session_id, reason)),
|
|
touch_media_frame=lambda session_id: touch_calls.append(session_id),
|
|
set_state=lambda session_id, state, handoff_reason, metadata: states.append((session_id, state)),
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
turns.append((session_id, transcript_text, barge_in))
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: handoffs.append(session_id),
|
|
handle_media_error=lambda session_id, message, metadata: errors.append((session_id, message)),
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
await asyncio.sleep(0.15)
|
|
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(speech_frame))
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(silence_frame))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
|
|
for _ in range(20):
|
|
if len(delivered) >= 2 and turns:
|
|
break
|
|
await asyncio.sleep(0.05)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert ("avs_media_runtime", "speaking") in states
|
|
assert ("avs_media_runtime", "thinking") in states
|
|
assert ("avs_media_runtime", "listening") in states
|
|
assert delivered[0] == ("avs_media_runtime", "greeting", True)
|
|
assert delivered[-1] == ("avs_media_runtime", "reply", False)
|
|
assert turns == [("avs_media_runtime", "hello", False)]
|
|
assert handoffs == []
|
|
assert touch_calls
|
|
assert media_ended
|
|
assert errors == []
|
|
|
|
|
|
def test_media_runtime_plays_filler_ack_when_v1_decision_is_slow():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_slow_v1",
|
|
call_id="call_media_runtime_slow_v1",
|
|
interaction_id="int_media_runtime_slow_v1",
|
|
ai_session_id="ais_media_runtime_slow_v1",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
|
|
def _slow_process_turn(session_id, transcript_text, language, barge_in, metadata):
|
|
del transcript_text, barge_in, metadata
|
|
time.sleep(0.9)
|
|
return VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime_slow_v1" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=_slow_process_turn,
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
await asyncio.sleep(0.15)
|
|
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(speech_frame))
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(silence_frame))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
|
|
for _ in range(60):
|
|
if len(delivered) >= 3:
|
|
break
|
|
await asyncio.sleep(0.05)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert delivered[0] == ("avs_media_runtime_slow_v1", "greeting", True)
|
|
assert delivered[1] == ("avs_media_runtime_slow_v1", "Секунду.", False)
|
|
assert delivered[2] == ("avs_media_runtime_slow_v1", "reply", False)
|
|
|
|
|
|
def test_media_runtime_speaks_technical_fallback_when_asr_transcribe_fails():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
errors: list[tuple[str, str]] = []
|
|
states: list[tuple[str, str]] = []
|
|
|
|
class _FailingASRProvider(ASRProvider):
|
|
name = "failing-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
del audio_bytes, language_hint
|
|
raise RuntimeError("AI_VOICE_ASR_YANDEX_API_KEY is required for Yandex ASR")
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_asr_error",
|
|
call_id="call_media_runtime_asr_error",
|
|
interaction_id="int_media_runtime_asr_error",
|
|
ai_session_id="ais_media_runtime_asr_error",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_FailingASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: states.append((session_id, state)),
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime_asr_error" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: errors.append((session_id, message)),
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
await asyncio.sleep(0.15)
|
|
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(speech_frame))
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(silence_frame))
|
|
await writer.drain()
|
|
|
|
for _ in range(40):
|
|
if len(delivered) >= 2 and errors:
|
|
break
|
|
await asyncio.sleep(0.05)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert delivered[0] == ("avs_media_runtime_asr_error", "greeting", True)
|
|
assert delivered[-1] == (
|
|
"avs_media_runtime_asr_error",
|
|
"Возникла техническая проблема со связью. Соединяю с оператором.",
|
|
False,
|
|
)
|
|
assert errors == [("avs_media_runtime_asr_error", "AI_VOICE_ASR_YANDEX_API_KEY is required for Yandex ASR")]
|
|
assert ("avs_media_runtime_asr_error", "handoff_requested") in states
|
|
|
|
|
|
def test_media_runtime_waits_for_late_registration():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
errors: list[tuple[str, str]] = []
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=1.0,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime_late" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: errors.append((session_id, message)),
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
async def _late_register() -> None:
|
|
await asyncio.sleep(0.15)
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_late",
|
|
call_id="call_media_runtime_late",
|
|
interaction_id="int_media_runtime_late",
|
|
ai_session_id="ais_media_runtime_late",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
|
|
late_task = asyncio.create_task(_late_register())
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
packet_type, _ = await read_packet(reader, timeout=2.0)
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
for _ in range(20):
|
|
if delivered:
|
|
break
|
|
await asyncio.sleep(0.05)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await late_task
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert delivered == [("avs_media_runtime_late", "greeting", True)]
|
|
assert errors == []
|
|
|
|
|
|
def test_media_runtime_sends_keepalive_while_tts_is_slow():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_keepalive",
|
|
call_id="call_media_runtime_keepalive",
|
|
interaction_id="int_media_runtime_keepalive",
|
|
ai_session_id="ais_media_runtime_keepalive",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
|
|
class _SlowTTSProvider(TTSProvider):
|
|
name = "slow-stub-tts"
|
|
|
|
def synthesize(
|
|
self,
|
|
text: str,
|
|
*,
|
|
language: str | None = None,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> TTSSynthesis:
|
|
del language, style_hints
|
|
assert text
|
|
time.sleep(1.2)
|
|
return TTSSynthesis(
|
|
text=text,
|
|
audio_bytes=(b"\x20\x00" * 960),
|
|
sample_rate_hz=24000,
|
|
)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_SlowTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime_keepalive" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
started_at = time.monotonic()
|
|
packet_type, payload = await read_packet(reader, timeout=1.0)
|
|
elapsed = time.monotonic() - started_at
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
assert payload == (b"\x00" * 320)
|
|
assert elapsed < 1.0
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
|
|
def test_media_runtime_sends_keepalive_before_registration_is_ready():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=2.0,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: "greeting" if session_id == "avs_media_runtime_prereg" else None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="answer",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
async def _late_register() -> None:
|
|
await asyncio.sleep(1.1)
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_prereg",
|
|
call_id="call_media_runtime_prereg",
|
|
interaction_id="int_media_runtime_prereg",
|
|
ai_session_id="ais_media_runtime_prereg",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
)
|
|
|
|
late_task = asyncio.create_task(_late_register())
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
started_at = time.monotonic()
|
|
packet_type, payload = await read_packet(reader, timeout=1.0)
|
|
elapsed = time.monotonic() - started_at
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
assert payload == (b"\x00" * 320)
|
|
assert elapsed < 1.0
|
|
|
|
for _ in range(30):
|
|
if delivered:
|
|
break
|
|
await asyncio.sleep(0.1)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await late_task
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert delivered == [("avs_media_runtime_prereg", "greeting", True)]
|
|
|
|
|
|
def test_media_runtime_starts_handoff_before_handoff_tts_finishes():
|
|
handoff_started = threading.Event()
|
|
events: list[str] = []
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="handoff_request",
|
|
reply_text="Соединяю с оператором",
|
|
confidence=0.9,
|
|
needs_handoff=True,
|
|
handoff_reason="Нужен живой оператор",
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="handoff requested",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="handoff_requested",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: (
|
|
events.append("handoff"),
|
|
handoff_started.set()
|
|
),
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
actor = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_handoff",
|
|
call_id="call_media_runtime_handoff",
|
|
interaction_id="int_media_runtime_handoff",
|
|
ai_session_id="ais_media_runtime_handoff",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
)
|
|
|
|
media_actor = None
|
|
|
|
async def _scenario() -> None:
|
|
nonlocal media_actor
|
|
media_actor = MediaActor(
|
|
registration=actor,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, text, is_greeting, style_hints
|
|
events.append("speak_start")
|
|
await asyncio.sleep(0)
|
|
assert handoff_started.wait(timeout=0.5)
|
|
await asyncio.sleep(0.05)
|
|
events.append("speak_end")
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(media_actor, pcm_frame, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert handoff_started.is_set()
|
|
assert events.index("handoff") < events.index("speak_end")
|
|
|
|
|
|
def test_media_runtime_voice_v2_plays_short_ack_before_main_reply():
|
|
planned: list[tuple[str, str, str, dict | None]] = []
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
|
|
class _ScheduleASRProvider(ASRProvider):
|
|
name = "schedule-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
del audio_bytes
|
|
return ASRTranscription(text="Хочу узнать график работы", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_ScheduleASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: planned.append((session_id, text, kind, metadata)),
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
time.sleep(0.25)
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите, пожалуйста, какой город вас интересует?",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
metadata={"early_intent": "schedule", "ack_kind": "understanding"},
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_write_audio_packet(current_actor, pcm_frame: bytes) -> None:
|
|
del current_actor, pcm_frame
|
|
|
|
runtime._write_audio_packet = _fake_write_audio_packet # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2",
|
|
call_id="call_media_runtime_v2",
|
|
interaction_id="int_media_runtime_v2",
|
|
ai_session_id="ais_media_runtime_v2",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert [item[2] for item in planned] == ["ack", "reply"]
|
|
assert planned[0][1] == "Сейчас сориентирую."
|
|
assert planned[1][1].startswith("Подскажите, пожалуйста")
|
|
assert delivered == [
|
|
("avs_media_runtime_v2", "Сейчас сориентирую.", False),
|
|
("avs_media_runtime_v2", "Подскажите, пожалуйста, какой город вас интересует?", False),
|
|
]
|
|
|
|
|
|
def test_media_runtime_voice_v2_uses_partial_asr_to_start_ack_before_full_asr():
|
|
timings: dict[str, float] = {}
|
|
speak_events: list[tuple[str, float]] = []
|
|
|
|
class _PartialAwareASRProvider(ASRProvider):
|
|
name = "partial-aware-asr"
|
|
|
|
def transcribe_partial(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
timings["partial_ready"] = time.monotonic()
|
|
return ASRTranscription(text="work schedule", language=language_hint or "ru", confidence=0.8)
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
timings["full_started"] = time.monotonic()
|
|
time.sleep(0.35)
|
|
timings["full_finished"] = time.monotonic()
|
|
return ASRTranscription(text="work schedule", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_PartialAwareASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите, какой город вас интересует?",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append((text, time.monotonic()))
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_partial",
|
|
call_id="call_media_runtime_v2_partial",
|
|
interaction_id="int_media_runtime_v2_partial",
|
|
ai_session_id="ais_media_runtime_v2_partial",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="listening",
|
|
)
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(35):
|
|
await runtime._handle_pcm(actor, speech_frame)
|
|
for _ in range(20):
|
|
if actor.partial_transcript:
|
|
break
|
|
await asyncio.sleep(0.01)
|
|
assert actor.partial_transcript == "work schedule"
|
|
for _ in range(2):
|
|
await runtime._handle_pcm(actor, silence_frame)
|
|
pcm_bytes, barge_in = await asyncio.wait_for(actor.turn_queue.get(), timeout=0.2)
|
|
await runtime._process_utterance(actor, pcm_bytes, barge_in)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert timings["partial_ready"] < timings["full_finished"]
|
|
assert speak_events
|
|
assert speak_events[0][0] == runtime._ack_text("ru", "understanding")
|
|
assert speak_events[0][1] < timings["full_finished"]
|
|
|
|
|
|
def test_media_runtime_voice_v2_emits_generic_ack_before_full_asr_without_partial_signal():
|
|
timings: dict[str, float] = {}
|
|
speak_events: list[tuple[str, float]] = []
|
|
|
|
class _SlowOnlyASRProvider(ASRProvider):
|
|
name = "slow-only-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
timings["full_started"] = time.monotonic()
|
|
time.sleep(0.35)
|
|
timings["full_finished"] = time.monotonic()
|
|
return ASRTranscription(text="hello there", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_SlowOnlyASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append((text, time.monotonic()))
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_generic_ack",
|
|
call_id="call_media_runtime_v2_generic_ack",
|
|
interaction_id="int_media_runtime_v2_generic_ack",
|
|
ai_session_id="ais_media_runtime_v2_generic_ack",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame * 45, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert speak_events
|
|
assert speak_events[0][0] == runtime._ack_text("ru", "generic")
|
|
assert speak_events[0][1] < timings["full_finished"]
|
|
|
|
|
|
def test_media_runtime_voice_v2_emits_ack_for_short_utterance_after_reduced_threshold():
|
|
timings: dict[str, float] = {}
|
|
speak_events: list[tuple[str, float]] = []
|
|
|
|
class _FastASRProvider(ASRProvider):
|
|
name = "fast-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
timings["full_started"] = time.monotonic()
|
|
time.sleep(0.12)
|
|
timings["full_finished"] = time.monotonic()
|
|
return ASRTranscription(text="hello there", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_FastASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append((text, time.monotonic()))
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_short_ack",
|
|
call_id="call_media_runtime_v2_short_ack",
|
|
interaction_id="int_media_runtime_v2_short_ack",
|
|
ai_session_id="ais_media_runtime_v2_short_ack",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert speak_events
|
|
assert speak_events[0][0] == runtime._ack_text("ru", "generic")
|
|
assert speak_events[0][1] < timings["full_finished"]
|
|
|
|
|
|
def test_media_runtime_voice_v2_emits_blind_ack_on_first_turn_without_partial_signal():
|
|
speak_events: list[tuple[str, float]] = []
|
|
|
|
class _FastASRProvider(ASRProvider):
|
|
name = "fast-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
time.sleep(0.08)
|
|
return ASRTranscription(text="hello there", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_FastASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append((text, time.monotonic()))
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_first_turn",
|
|
call_id="call_media_runtime_v2_first_turn",
|
|
interaction_id="int_media_runtime_v2_first_turn",
|
|
ai_session_id="ais_media_runtime_v2_first_turn",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert len(speak_events) == 2
|
|
assert speak_events[0][0] == runtime._ack_text("ru", "generic")
|
|
assert speak_events[1][0] == "Подскажите подробнее, пожалуйста."
|
|
|
|
|
|
def test_media_runtime_voice_v2_skips_filler_ack_when_caller_says_goodbye():
|
|
planned: list[tuple[str, str, str, dict | None]] = []
|
|
|
|
class _GoodbyeASRProvider(ASRProvider):
|
|
name = "goodbye-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="Спасибо, до свидания", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_GoodbyeASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: planned.append((session_id, text, kind, metadata)),
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
time.sleep(0.25)
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="closing",
|
|
reply_text="Хорошо, всего доброго!",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="close",
|
|
kb_refs=[],
|
|
summary_text="call wrapped up",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_goodbye",
|
|
call_id="call_media_runtime_v2_goodbye",
|
|
interaction_id="int_media_runtime_v2_goodbye",
|
|
ai_session_id="ais_media_runtime_v2_goodbye",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert [item[2] for item in planned] == ["reply"]
|
|
assert planned[0][1] == "Хорошо, всего доброго!"
|
|
|
|
|
|
def test_media_runtime_voice_v2_throttles_repeated_filler_ack_within_gap():
|
|
planned: list[tuple[str, str, str, dict | None]] = []
|
|
|
|
class _SlowASRProvider(ASRProvider):
|
|
name = "slow-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="Хочу узнать график работы", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_SlowASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: planned.append((session_id, text, kind, metadata)),
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
time.sleep(0.25)
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите, пожалуйста, какой город вас интересует?",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_throttle",
|
|
call_id="call_media_runtime_v2_throttle",
|
|
interaction_id="int_media_runtime_v2_throttle",
|
|
ai_session_id="ais_media_runtime_v2_throttle",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
runtime._reset_live_turn_state(actor)
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert [item[2] for item in planned] == ["ack", "reply", "reply"]
|
|
|
|
|
|
def test_media_runtime_voice_v2_blind_ack_defers_to_known_low_signal_partial_transcript():
|
|
speak_events: list[str] = []
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append(text)
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_low_signal_blind",
|
|
call_id="call_media_runtime_v2_low_signal_blind",
|
|
interaction_id="int_media_runtime_v2_low_signal_blind",
|
|
ai_session_id="ais_media_runtime_v2_low_signal_blind",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
# A caller who already said a recognized filler-answer ("да") should not
|
|
# get a blind ack just because the audio clip crossed the length threshold.
|
|
actor.stable_partial_transcript = "да"
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert speak_events == ["Подскажите подробнее, пожалуйста."]
|
|
|
|
|
|
def test_media_runtime_low_signal_filter_catches_short_asr_noise():
|
|
assert AudioSocketMediaRuntime._is_low_signal_partial_transcript("Давай")
|
|
assert AudioSocketMediaRuntime._is_low_signal_partial_transcript("твой")
|
|
assert AudioSocketMediaRuntime._is_low_signal_partial_transcript("Поргай, что это")
|
|
assert AudioSocketMediaRuntime._is_low_signal_final_transcript("твой")
|
|
assert AudioSocketMediaRuntime._is_low_signal_final_transcript("Поргай, что это")
|
|
|
|
|
|
def test_media_runtime_final_low_signal_filter_keeps_real_answers():
|
|
# A finalized "да"/"нет"/etc. is a real answer, not noise — must not be silently dropped.
|
|
for real_answer in ("Да", "Нет", "Хорошо", "Ладно", "Давай", "Привет"):
|
|
assert not AudioSocketMediaRuntime._is_low_signal_final_transcript(real_answer)
|
|
|
|
|
|
def test_media_runtime_ignores_low_signal_utterance_without_ack_or_turn():
|
|
planned: list[tuple[str, str, str, dict | None]] = []
|
|
delivered: list[tuple[str, str, bool]] = []
|
|
turns: list[str] = []
|
|
|
|
class _LowSignalASRProvider(ASRProvider):
|
|
name = "low-signal-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="Ой", language=language_hint or "ru", confidence=0.82)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_LowSignalASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((session_id, text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: planned.append((session_id, text, kind, metadata)),
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
turns.append(transcript_text)
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_low_signal",
|
|
call_id="call_media_runtime_low_signal",
|
|
interaction_id="int_media_runtime_low_signal",
|
|
ai_session_id="ais_media_runtime_low_signal",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame * 12, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert turns == []
|
|
assert planned == []
|
|
assert delivered == []
|
|
|
|
|
|
def test_media_runtime_voice_v2_inserts_small_gap_between_ack_and_main_reply():
|
|
speak_events: list[tuple[str, float]] = []
|
|
|
|
class _InstantASRProvider(ASRProvider):
|
|
name = "instant-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="hello there", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_InstantASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_text(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints
|
|
speak_events.append((text, time.monotonic()))
|
|
await asyncio.sleep(0)
|
|
|
|
runtime._speak_text = _fake_speak_text # type: ignore[method-assign]
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_gap",
|
|
call_id="call_media_runtime_v2_gap",
|
|
interaction_id="int_media_runtime_v2_gap",
|
|
ai_session_id="ais_media_runtime_v2_gap",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=False,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert len(speak_events) == 2
|
|
assert speak_events[0][0] == runtime._ack_text("ru", "generic")
|
|
assert speak_events[1][0] == "Подскажите подробнее, пожалуйста."
|
|
assert (speak_events[1][1] - speak_events[0][1]) >= (runtime._v2_ack_post_gap_seconds - 0.02)
|
|
|
|
|
|
def test_media_runtime_voice_v2_emotive_ack_uses_ru_variants_and_style_hints_only_for_ack():
|
|
synth_calls: list[tuple[str, dict[str, object] | None]] = []
|
|
|
|
class _RecordingTTSProvider(TTSProvider):
|
|
name = "recording-tts"
|
|
|
|
def synthesize(
|
|
self,
|
|
text: str,
|
|
*,
|
|
language: str | None = None,
|
|
style_hints: dict[str, object] | None = None,
|
|
) -> TTSSynthesis:
|
|
del language
|
|
synth_calls.append((text, style_hints))
|
|
return TTSSynthesis(text=text, audio_bytes=(b"\x10\x00" * 960), sample_rate_hz=24000)
|
|
|
|
class _ScheduleASRProvider(ASRProvider):
|
|
name = "schedule-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
del audio_bytes
|
|
return ASRTranscription(text="Хочу узнать график работы", language=language_hint or "ru", confidence=0.9)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_ScheduleASRProvider(),
|
|
tts_provider=_RecordingTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
time.sleep(0.25)
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите, пожалуйста, какой город вас интересует?",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
metadata={"early_intent": "schedule", "ack_kind": "understanding"},
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_write_audio_packet(current_actor, pcm_frame: bytes) -> None:
|
|
del current_actor, pcm_frame
|
|
|
|
runtime._write_audio_packet = _fake_write_audio_packet # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_emotive_ack",
|
|
call_id="call_media_runtime_v2_emotive_ack",
|
|
interaction_id="int_media_runtime_v2_emotive_ack",
|
|
ai_session_id="ais_media_runtime_v2_emotive_ack",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=False,
|
|
voice_v2_partial_asr=False,
|
|
voice_v2_emotive_ack=True,
|
|
voice_v2_emotive_ack_ru_only=True,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.finalized_caller_turn_count = 1
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert len(synth_calls) == 2
|
|
assert synth_calls[0][0] in runtime._base_ack_variants("ru", "understanding")
|
|
assert synth_calls[0][1] == {"role": "good"}
|
|
assert synth_calls[1][0] == "Подскажите, пожалуйста, какой город вас интересует?"
|
|
assert synth_calls[1][1] is None
|
|
|
|
|
|
def test_media_runtime_voice_v2_emotive_ack_avoids_same_variant_back_to_back():
|
|
async def _scenario() -> tuple[str, str, AudioSocketMediaRuntime]:
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_repeat_guard",
|
|
call_id="call_media_runtime_v2_repeat_guard",
|
|
interaction_id="int_media_runtime_v2_repeat_guard",
|
|
ai_session_id="ais_media_runtime_v2_repeat_guard",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_emotive_ack=True,
|
|
voice_v2_emotive_ack_ru_only=True,
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.response_plan_id = "rsp_same_seed"
|
|
|
|
first_text, _, _ = runtime._select_ack_payload(actor, language="ru", ack_kind="understanding")
|
|
actor.response_plan_id = "rsp_same_seed"
|
|
second_text, _, _ = runtime._select_ack_payload(actor, language="ru", ack_kind="understanding")
|
|
return first_text, second_text, runtime
|
|
|
|
first_text, second_text, runtime = asyncio.run(_scenario())
|
|
|
|
assert first_text in runtime._base_ack_variants("ru", "understanding")
|
|
assert second_text in runtime._base_ack_variants("ru", "understanding")
|
|
assert first_text != second_text
|
|
|
|
|
|
def test_media_runtime_streaming_timeout_enters_backoff_before_reopen():
|
|
class _FlakyStreamingProvider(StreamingASRProvider):
|
|
name = "flaky-streaming"
|
|
supports_streaming = True
|
|
|
|
def __init__(self) -> None:
|
|
self.open_calls = 0
|
|
|
|
def open_stream(self, session_id: str, *, language_hint: str | None = None) -> str:
|
|
del session_id, language_hint
|
|
self.open_calls += 1
|
|
return "stream-1"
|
|
|
|
def push_pcm(self, stream_id: str, pcm_8k_chunk: bytes) -> None:
|
|
del stream_id, pcm_8k_chunk
|
|
raise StreamingASRUnavailable("timed out")
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_stream_backoff",
|
|
call_id="call_media_runtime_stream_backoff",
|
|
interaction_id="int_media_runtime_stream_backoff",
|
|
ai_session_id="ais_media_runtime_stream_backoff",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
)
|
|
streaming_provider = _FlakyStreamingProvider()
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=streaming_provider,
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="listening",
|
|
)
|
|
runtime._reset_live_turn_state(actor)
|
|
await runtime._ensure_streaming_asr(actor)
|
|
assert actor.asr_streaming_enabled is True
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._handle_pcm(actor, speech_frame)
|
|
await runtime._handle_pcm(actor, speech_frame)
|
|
for _ in range(20):
|
|
if not actor.asr_streaming_enabled:
|
|
break
|
|
await asyncio.sleep(0.01)
|
|
assert actor.asr_streaming_enabled is False
|
|
assert actor.streaming_asr_backoff_until_monotonic > time.monotonic()
|
|
await runtime._ensure_streaming_asr(actor)
|
|
assert streaming_provider.open_calls == 1
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
|
|
def test_media_runtime_streaming_sidecar_push_does_not_block_vad_finalization():
|
|
class _SlowStreamingProvider(StreamingASRProvider):
|
|
name = "slow-streaming"
|
|
supports_streaming = True
|
|
|
|
def __init__(self) -> None:
|
|
self.push_count = 0
|
|
|
|
def open_stream(self, session_id: str, *, language_hint: str | None = None) -> str:
|
|
del session_id, language_hint
|
|
return "stream-1"
|
|
|
|
def push_pcm(self, stream_id: str, pcm_8k_chunk: bytes) -> None:
|
|
del stream_id, pcm_8k_chunk
|
|
self.push_count += 1
|
|
time.sleep(0.25)
|
|
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_nonblocking_asr",
|
|
call_id="call_media_runtime_nonblocking_asr",
|
|
interaction_id="int_media_runtime_nonblocking_asr",
|
|
ai_session_id="ais_media_runtime_nonblocking_asr",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
)
|
|
streaming_provider = _SlowStreamingProvider()
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=streaming_provider,
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="reply",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> float:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="listening",
|
|
)
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
started = time.monotonic()
|
|
await runtime._handle_pcm(actor, speech_frame)
|
|
await runtime._handle_pcm(actor, speech_frame)
|
|
await runtime._handle_pcm(actor, silence_frame)
|
|
await runtime._handle_pcm(actor, silence_frame)
|
|
elapsed = time.monotonic() - started
|
|
pcm_bytes, _ = await asyncio.wait_for(actor.turn_queue.get(), timeout=0.1)
|
|
assert pcm_bytes
|
|
for _ in range(50):
|
|
if streaming_provider.push_count >= 1:
|
|
break
|
|
await asyncio.sleep(0.01)
|
|
await runtime._close_streaming_asr(actor, drain=False)
|
|
return elapsed
|
|
|
|
elapsed = asyncio.run(_scenario())
|
|
|
|
assert elapsed < 0.15
|
|
assert streaming_provider.push_count >= 1
|
|
|
|
|
|
def test_media_runtime_streaming_partial_poll_does_not_block_turn_close_ack():
|
|
events: list[str] = []
|
|
|
|
class _StreamingProvider(StreamingASRProvider):
|
|
name = "streaming-sidecar"
|
|
supports_streaming = True
|
|
|
|
def __init__(self) -> None:
|
|
self.poll_count = 0
|
|
|
|
def poll_partial(self, stream_id: str) -> StreamingASRPartial | None:
|
|
del stream_id
|
|
self.poll_count += 1
|
|
time.sleep(0.25)
|
|
return StreamingASRPartial(text="need schedule", language="ru", confidence=0.8, is_stable=True)
|
|
|
|
def finalize(self, stream_id: str) -> ASRTranscription:
|
|
assert stream_id == "stream-1"
|
|
events.append("finalize")
|
|
return ASRTranscription(text="need schedule", language="ru", confidence=0.9)
|
|
|
|
def close_stream(self, stream_id: str) -> None:
|
|
assert stream_id == "stream-1"
|
|
|
|
streaming_provider = _StreamingProvider()
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=streaming_provider,
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="schedule",
|
|
reply_text="reply",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_reply(
|
|
actor: MediaActor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
reply_phase: str | None = "main",
|
|
) -> None:
|
|
del actor, text, is_greeting, style_hints
|
|
events.append(str(reply_phase))
|
|
|
|
runtime._speak_reply = _fake_speak_reply # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_turn_close_no_poll",
|
|
call_id="call_media_runtime_turn_close_no_poll",
|
|
interaction_id="int_media_runtime_turn_close_no_poll",
|
|
ai_session_id="ais_media_runtime_turn_close_no_poll",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.asr_streaming_enabled = True
|
|
actor.asr_stream_id = "stream-1"
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert streaming_provider.poll_count == 0
|
|
assert events[:2] == ["ack", "finalize"]
|
|
assert "main" in events
|
|
|
|
|
|
def test_media_runtime_voice_v2_uses_early_plan_before_slow_final_decision():
|
|
events: list[str | tuple[str, str]] = []
|
|
|
|
class _StreamingProvider(StreamingASRProvider):
|
|
name = "streaming-sidecar"
|
|
supports_streaming = True
|
|
|
|
def finalize(self, stream_id: str) -> ASRTranscription:
|
|
assert stream_id == "stream-1"
|
|
return ASRTranscription(text="work schedule", language="ru", confidence=0.9)
|
|
|
|
def close_stream(self, stream_id: str) -> None:
|
|
assert stream_id == "stream-1"
|
|
|
|
def _process_turn(session_id, transcript_text, language, barge_in, metadata):
|
|
del session_id, transcript_text, barge_in
|
|
if metadata and metadata.get("reply_phase") == "early_plan":
|
|
events.append("early_plan")
|
|
return VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="schedule",
|
|
reply_text="early reply",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="early ready",
|
|
model="early",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
events.append("final_start")
|
|
time.sleep(0.25)
|
|
events.append("final_done")
|
|
return VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="schedule",
|
|
reply_text="final reply",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="final ready",
|
|
model="final",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=_StreamingProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=_process_turn,
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_reply(
|
|
actor: MediaActor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
reply_phase: str | None = "main",
|
|
) -> None:
|
|
del actor, is_greeting, style_hints
|
|
events.append((str(reply_phase), text))
|
|
|
|
runtime._speak_reply = _fake_speak_reply # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_early_plan",
|
|
call_id="call_media_runtime_early_plan",
|
|
interaction_id="int_media_runtime_early_plan",
|
|
ai_session_id="ais_media_runtime_early_plan",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.asr_streaming_enabled = True
|
|
actor.asr_stream_id = "stream-1"
|
|
actor.partial_transcript = "work schedule"
|
|
actor.stable_partial_transcript = "work schedule"
|
|
actor.partial_intent = "schedule"
|
|
actor.stable_partial_intent = "schedule"
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
for _ in range(50):
|
|
if "final_done" in events:
|
|
break
|
|
await asyncio.sleep(0.01)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert ("main", "early reply") in events
|
|
assert ("main", "final reply") not in events
|
|
assert events.index(("main", "early reply")) < events.index("final_done")
|
|
|
|
|
|
def test_media_runtime_voice_v2_uses_partial_as_final_when_streaming_finalize_fails():
|
|
captured: list[tuple[str, dict | None]] = []
|
|
|
|
class _CountingASRProvider(ASRProvider):
|
|
name = "counting-asr"
|
|
|
|
def __init__(self) -> None:
|
|
self.transcribe_count = 0
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
del audio_bytes
|
|
self.transcribe_count += 1
|
|
return ASRTranscription(text="batch text", language=language_hint or "ru", confidence=0.9)
|
|
|
|
class _FailingStreamingProvider(StreamingASRProvider):
|
|
name = "failing-streaming"
|
|
supports_streaming = True
|
|
|
|
def finalize(self, stream_id: str) -> ASRTranscription:
|
|
assert stream_id == "stream-1"
|
|
raise StreamingASRUnavailable("finalize timeout")
|
|
|
|
def close_stream(self, stream_id: str) -> None:
|
|
assert stream_id == "stream-1"
|
|
|
|
asr_provider = _CountingASRProvider()
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=asr_provider,
|
|
streaming_asr_provider=_FailingStreamingProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
captured.append((transcript_text, metadata))
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="schedule",
|
|
reply_text="reply",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _fake_speak_reply(
|
|
actor: MediaActor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
reply_phase: str | None = "main",
|
|
) -> None:
|
|
del actor, text, is_greeting, style_hints, reply_phase
|
|
|
|
runtime._speak_reply = _fake_speak_reply # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_partial_final",
|
|
call_id="call_media_runtime_partial_final",
|
|
interaction_id="int_media_runtime_partial_final",
|
|
ai_session_id="ais_media_runtime_partial_final",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
actor.asr_streaming_enabled = True
|
|
actor.asr_stream_id = "stream-1"
|
|
actor.partial_transcript = "work schedule"
|
|
actor.stable_partial_transcript = "work schedule"
|
|
actor.partial_intent = "schedule"
|
|
actor.stable_partial_intent = "schedule"
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame * 40, False)
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert asr_provider.transcribe_count == 0
|
|
final_turns = [item for item in captured if item[1] and item[1].get("reply_phase") == "final"]
|
|
assert final_turns
|
|
assert final_turns[0][0] == "work schedule"
|
|
assert final_turns[0][1]["transcript_source"] == "streaming_partial_after_finalize_failure"
|
|
|
|
|
|
def test_media_runtime_merges_thinking_continuation_into_current_utterance():
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=StreamingASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
base_pcm = (1000).to_bytes(2, "little", signed=True) * 160
|
|
continuation_pcm = (900).to_bytes(2, "little", signed=True) * 160
|
|
|
|
async def _scenario() -> bytes:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_thinking_merge",
|
|
call_id="call_media_runtime_thinking_merge",
|
|
interaction_id="int_media_runtime_thinking_merge",
|
|
ai_session_id="ais_media_runtime_thinking_merge",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="thinking",
|
|
)
|
|
await runtime._handle_pcm(actor, continuation_pcm)
|
|
await runtime._handle_pcm(actor, continuation_pcm)
|
|
merged = await runtime._extend_with_thinking_continuation(actor, base_pcm)
|
|
assert actor.thinking_continuation_pcm == bytearray()
|
|
return merged
|
|
|
|
merged = asyncio.run(_scenario())
|
|
|
|
assert len(merged) > len(base_pcm)
|
|
assert merged.endswith(continuation_pcm + continuation_pcm)
|
|
|
|
|
|
def test_media_runtime_thinking_continuation_is_not_kept_alive_by_silence():
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=StreamingASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
base_pcm = (1000).to_bytes(2, "little", signed=True) * 160
|
|
continuation_pcm = (900).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
|
|
async def _scenario() -> bytes:
|
|
actor = MediaActor(
|
|
registration=MediaRegistration(
|
|
voice_session_id="avs_media_runtime_thinking_silence_tail",
|
|
call_id="call_media_runtime_thinking_silence_tail",
|
|
interaction_id="int_media_runtime_thinking_silence_tail",
|
|
ai_session_id="ais_media_runtime_thinking_silence_tail",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
),
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="thinking",
|
|
)
|
|
merge_task = asyncio.create_task(runtime._extend_with_thinking_continuation(actor, base_pcm))
|
|
await runtime._handle_pcm(actor, continuation_pcm)
|
|
for _ in range(40):
|
|
await runtime._handle_pcm(actor, silence_frame)
|
|
await asyncio.sleep(0.02)
|
|
merged = await asyncio.wait_for(merge_task, timeout=0.2)
|
|
assert actor.thinking_continuation_pcm == bytearray()
|
|
return merged
|
|
|
|
merged = asyncio.run(_scenario())
|
|
|
|
assert len(merged) > len(base_pcm)
|
|
assert continuation_pcm in merged
|
|
|
|
|
|
def _legacy_test_media_runtime_voice_v2_uses_streaming_sidecar_for_partial_and_final_asr():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
reply_starts: list[tuple[str, str | None, float]] = []
|
|
turns: list[str] = []
|
|
turn_ready = threading.Event()
|
|
finalize_ready = threading.Event()
|
|
replies_ready = threading.Event()
|
|
delivered: list[tuple[str, bool]] = []
|
|
|
|
class _ExplodingBatchASRProvider(ASRProvider):
|
|
name = "exploding-batch-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
raise AssertionError("batch ASR should not be used when streaming sidecar is active")
|
|
|
|
class _FakeStreamingASRProvider(StreamingASRProvider):
|
|
name = "fake-streaming-sidecar"
|
|
supports_streaming = True
|
|
|
|
def __init__(self) -> None:
|
|
self.chunk_count = 0
|
|
self.events: list[tuple[str, float]] = []
|
|
|
|
def open_stream(self, session_id: str, *, language_hint: str | None = None) -> str:
|
|
del session_id, language_hint
|
|
self.events.append(("open", time.monotonic()))
|
|
return "stream-1"
|
|
|
|
def push_pcm(self, stream_id: str, pcm_8k_chunk: bytes) -> None:
|
|
del stream_id
|
|
assert pcm_8k_chunk
|
|
self.chunk_count += 1
|
|
self.events.append(("push", time.monotonic()))
|
|
|
|
def poll_partial(self, stream_id: str) -> StreamingASRPartial | None:
|
|
del stream_id
|
|
self.events.append(("poll", time.monotonic()))
|
|
if self.chunk_count >= 2:
|
|
return StreamingASRPartial(
|
|
text="мне нужен график работы",
|
|
language="ru",
|
|
confidence=0.84,
|
|
is_stable=True,
|
|
)
|
|
return None
|
|
|
|
def finalize(self, stream_id: str) -> ASRTranscription:
|
|
del stream_id
|
|
self.events.append(("finalize", time.monotonic()))
|
|
finalize_ready.set()
|
|
time.sleep(0.2)
|
|
return ASRTranscription(text="мне нужен график работы", language="ru", confidence=0.84)
|
|
|
|
def close_stream(self, stream_id: str) -> None:
|
|
del stream_id
|
|
self.events.append(("close", time.monotonic()))
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_streaming",
|
|
call_id="call_media_runtime_v2_streaming",
|
|
interaction_id="int_media_runtime_v2_streaming",
|
|
ai_session_id="ais_media_runtime_v2_streaming",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
)
|
|
streaming_provider = _FakeStreamingASRProvider()
|
|
|
|
def _mark_reply_started(session_id, text, is_greeting, phase):
|
|
del session_id, is_greeting
|
|
reply_starts.append((text, phase, time.monotonic()))
|
|
if len(reply_starts) >= 2:
|
|
replies_ready.set()
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_ExplodingBatchASRProvider(),
|
|
streaming_asr_provider=streaming_provider,
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_started=_mark_reply_started,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: delivered.append((text, is_greeting)),
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
turns.append(transcript_text)
|
|
or turn_ready.set()
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Назовите, пожалуйста, город.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(14):
|
|
writer.write(encode_audio_packet(speech_frame))
|
|
for _ in range(2):
|
|
writer.write(encode_audio_packet(silence_frame))
|
|
await writer.drain()
|
|
|
|
assert await asyncio.to_thread(finalize_ready.wait, 10.0)
|
|
assert await asyncio.to_thread(turn_ready.wait, 10.0)
|
|
assert await asyncio.to_thread(replies_ready.wait, 10.0)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
phases = [phase for _, phase, _ in reply_starts]
|
|
finalize_started_at = next(ts for name, ts in streaming_provider.events if name == "finalize")
|
|
assert turns == ["мне нужен график работы"]
|
|
assert "open" in [name for name, _ in streaming_provider.events]
|
|
assert "close" in [name for name, _ in streaming_provider.events]
|
|
assert phases[:2] == ["ack", "main"]
|
|
assert reply_starts[0][2] < finalize_started_at
|
|
assert any(text == "Назовите, пожалуйста, город." and phase == "main" for text, phase, _ in reply_starts)
|
|
|
|
|
|
def _legacy_test_media_runtime_voice_v2_falls_back_when_streaming_sidecar_is_unavailable():
|
|
registrations: dict[str, MediaRegistration] = {}
|
|
turns: list[str] = []
|
|
turn_ready = threading.Event()
|
|
|
|
class _BatchASRProvider(ASRProvider):
|
|
name = "batch-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="нужен оператор", language=language_hint or "ru", confidence=0.88)
|
|
|
|
class _UnavailableStreamingProvider(StreamingASRProvider):
|
|
name = "missing-sidecar"
|
|
supports_streaming = True
|
|
|
|
def open_stream(self, session_id: str, *, language_hint: str | None = None) -> str:
|
|
del session_id, language_hint
|
|
raise StreamingASRUnavailable("sidecar down")
|
|
|
|
media_uuid = str(uuid.uuid4())
|
|
registrations[media_uuid] = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_v2_fallback",
|
|
call_id="call_media_runtime_v2_fallback",
|
|
interaction_id="int_media_runtime_v2_fallback",
|
|
ai_session_id="ais_media_runtime_v2_fallback",
|
|
language="ru",
|
|
media_uuid=media_uuid,
|
|
queue_code="voice_lab_ai",
|
|
queue_id="que_voice_lab_ai",
|
|
agent_profile="voice_support",
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
)
|
|
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=2000,
|
|
asr_provider=_BatchASRProvider(),
|
|
streaming_asr_provider=_UnavailableStreamingProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: registrations.get(value),
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: (
|
|
turns.append(transcript_text)
|
|
or turn_ready.set()
|
|
or VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="handoff",
|
|
reply_text="Соединяю с оператором.",
|
|
confidence=0.9,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
)
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
await runtime.start()
|
|
port = runtime._server.sockets[0].getsockname()[1]
|
|
|
|
reader, writer = await asyncio.open_connection("127.0.0.1", port)
|
|
writer.write(encode_packet(AUDIO_SOCKET_PACKET_UUID, uuid.UUID(media_uuid).bytes))
|
|
await writer.drain()
|
|
|
|
speech_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
silence_frame = b"\x00\x00" * 160
|
|
for _ in range(4):
|
|
writer.write(encode_audio_packet(speech_frame))
|
|
for _ in range(4):
|
|
writer.write(encode_audio_packet(silence_frame))
|
|
await writer.drain()
|
|
|
|
assert await asyncio.to_thread(turn_ready.wait, 10.0)
|
|
|
|
writer.close()
|
|
await writer.wait_closed()
|
|
await asyncio.sleep(0.2)
|
|
await runtime.stop()
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert turns == ["нужен оператор"]
|
|
assert registrations[media_uuid].voice_v2_duplex is False
|
|
assert registrations[media_uuid].voice_v2_partial_asr is False
|
|
|
|
|
|
def test_media_runtime_forwards_vad_preroll_to_streaming_asr_on_speech_start():
|
|
class _RecordingStreamingProvider(StreamingASRProvider):
|
|
name = "recording-streaming"
|
|
supports_streaming = True
|
|
|
|
def __init__(self) -> None:
|
|
self.pushed_chunks: list[bytes] = []
|
|
|
|
def open_stream(self, session_id: str, *, language_hint: str | None = None) -> str:
|
|
del session_id, language_hint
|
|
return "stream-1"
|
|
|
|
def push_pcm(self, stream_id: str, pcm_8k_chunk: bytes) -> None:
|
|
assert stream_id == "stream-1"
|
|
self.pushed_chunks.append(pcm_8k_chunk)
|
|
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_preroll",
|
|
call_id="call_media_runtime_preroll",
|
|
interaction_id="int_media_runtime_preroll",
|
|
ai_session_id="ais_media_runtime_preroll",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
voice_v2_enabled=True,
|
|
voice_v2_ack_mode="immediate_short",
|
|
voice_v2_streaming_tts=True,
|
|
voice_v2_partial_asr=True,
|
|
voice_v2_duplex=True,
|
|
voice_v2_streaming_asr_backend="local_sidecar",
|
|
)
|
|
streaming_provider = _RecordingStreamingProvider()
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=1000,
|
|
max_turn_ms=2000,
|
|
asr_provider=_StubASRProvider(),
|
|
streaming_asr_provider=streaming_provider,
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="clarification",
|
|
reply_text="Подскажите подробнее, пожалуйста.",
|
|
confidence=0.8,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="reply ready",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=1000, max_turn_ms=2000),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="listening",
|
|
)
|
|
first_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
second_frame = (1200).to_bytes(2, "little", signed=True) * 160
|
|
# min_speech_ms=40 with frame_ms=20 means speech_started only fires on
|
|
# the *second* speech frame; the first frame is buffered as VAD pre-roll.
|
|
await runtime._handle_pcm(actor, first_frame)
|
|
assert streaming_provider.pushed_chunks == []
|
|
await runtime._handle_pcm(actor, second_frame)
|
|
assert actor.asr_streaming_enabled is True
|
|
queue = actor.streaming_asr_push_queue
|
|
assert queue is not None
|
|
await asyncio.wait_for(queue.join(), timeout=1.0)
|
|
pushed = b"".join(streaming_provider.pushed_chunks)
|
|
assert first_frame in pushed, "VAD pre-roll audio must reach the streaming ASR stream"
|
|
assert second_frame in pushed
|
|
assert pushed.count(first_frame) == 1, "pre-roll must not be duplicated"
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
|
|
def test_media_runtime_reprompts_after_consecutive_low_signal_turns():
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_reprompt_low_signal",
|
|
call_id="call_media_runtime_reprompt_low_signal",
|
|
interaction_id="int_media_runtime_reprompt_low_signal",
|
|
ai_session_id="ais_media_runtime_reprompt_low_signal",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
)
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: VoiceAITurnDecisionOut(
|
|
language=language or "ru",
|
|
intent="unknown",
|
|
reply_text="",
|
|
confidence=0.1,
|
|
needs_handoff=False,
|
|
handoff_reason=None,
|
|
case_action="keep_open",
|
|
kb_refs=[],
|
|
summary_text="noise",
|
|
model="stub-voice",
|
|
latency_ms=1,
|
|
status="active",
|
|
),
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
runtime._no_speech_reprompt_low_signal_turns = 2
|
|
|
|
class _LowSignalASRProvider(ASRProvider):
|
|
name = "low-signal-asr"
|
|
|
|
def transcribe(self, audio_bytes: bytes, *, language_hint: str | None = None) -> ASRTranscription:
|
|
assert audio_bytes
|
|
return ASRTranscription(text="алло", language=language_hint or "ru", confidence=0.2)
|
|
|
|
runtime._asr_provider = _LowSignalASRProvider()
|
|
|
|
spoken: list[str] = []
|
|
|
|
async def _fake_speak_reply(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
reply_phase: str | None = "main",
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints, reply_phase
|
|
spoken.append(text)
|
|
|
|
runtime._speak_reply = _fake_speak_reply # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
state="listening",
|
|
)
|
|
pcm_frame = (1000).to_bytes(2, "little", signed=True) * 160
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
assert spoken == [], "should not reprompt after a single low-signal turn"
|
|
assert actor.consecutive_low_signal_turns == 1
|
|
|
|
await runtime._process_utterance(actor, pcm_frame, False)
|
|
assert len(spoken) == 1, "should reprompt after reaching the low-signal turn threshold"
|
|
assert actor.consecutive_low_signal_turns == 0
|
|
assert actor.reprompt_attempts == 1
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
|
|
def test_media_runtime_no_speech_watchdog_reprompts_on_silence_timeout():
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_reprompt_silence",
|
|
call_id="call_media_runtime_reprompt_silence",
|
|
interaction_id="int_media_runtime_reprompt_silence",
|
|
ai_session_id="ais_media_runtime_reprompt_silence",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
)
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_StubTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: None,
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
runtime._no_speech_reprompt_seconds = 0.05
|
|
|
|
spoken: list[str] = []
|
|
|
|
async def _fake_speak_reply(
|
|
current_actor,
|
|
text: str,
|
|
*,
|
|
is_greeting: bool,
|
|
style_hints: dict[str, object] | None = None,
|
|
reply_phase: str | None = "main",
|
|
) -> None:
|
|
del current_actor, is_greeting, style_hints, reply_phase
|
|
spoken.append(text)
|
|
|
|
runtime._speak_reply = _fake_speak_reply # type: ignore[method-assign]
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=None, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
# The no-speech watchdog check now rides the existing keepalive loop
|
|
# instead of its own dedicated task; keep last_outbound_audio_monotonic
|
|
# fresh so the loop's silence-keepalive-audio branch (which needs a
|
|
# real writer) doesn't fire during this test.
|
|
actor.last_outbound_audio_monotonic = time.monotonic()
|
|
await runtime._set_actor_state(actor, "listening")
|
|
keepalive_task = asyncio.create_task(runtime._keepalive_loop(actor))
|
|
for _ in range(20):
|
|
if spoken:
|
|
break
|
|
await asyncio.sleep(0.05)
|
|
actor.closed = True
|
|
keepalive_task.cancel()
|
|
with contextlib.suppress(asyncio.CancelledError, Exception):
|
|
await keepalive_task
|
|
assert spoken, "watchdog should reprompt after prolonged silence in listening state"
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
|
|
def test_media_runtime_streaming_tts_prebuffers_without_dropping_audio():
|
|
class _ChunkedTTSProvider(TTSProvider):
|
|
name = "chunked-tts"
|
|
|
|
def synthesize(self, text, *, language=None, style_hints=None):
|
|
raise AssertionError("streaming session should use synthesize_chunks")
|
|
|
|
def synthesize_chunks(self, text, *, language=None, style_hints=None):
|
|
del language, style_hints
|
|
# 12 small 16kHz chunks (5ms each), each well under the prebuffer
|
|
# target, so both the accumulation path and the trailing flush
|
|
# (for whatever is still buffered once the stream ends) get exercised.
|
|
tone = (500).to_bytes(2, "little", signed=True) * 80
|
|
for _ in range(12):
|
|
yield TTSSynthesis(text=text, audio_bytes=tone, sample_rate_hz=16000)
|
|
|
|
class _FakeWriter:
|
|
def __init__(self) -> None:
|
|
self.packets: list[bytes] = []
|
|
|
|
def write(self, data: bytes) -> None:
|
|
self.packets.append(data)
|
|
|
|
async def drain(self) -> None:
|
|
return None
|
|
|
|
registration = MediaRegistration(
|
|
voice_session_id="avs_media_runtime_prebuffer",
|
|
call_id="call_media_runtime_prebuffer",
|
|
interaction_id="int_media_runtime_prebuffer",
|
|
ai_session_id="ais_media_runtime_prebuffer",
|
|
language="ru",
|
|
media_uuid=str(uuid.uuid4()),
|
|
voice_v2_streaming_tts=True,
|
|
)
|
|
runtime = AudioSocketMediaRuntime(
|
|
enabled=True,
|
|
host="127.0.0.1",
|
|
port=0,
|
|
frame_ms=20,
|
|
idle_timeout_seconds=2.0,
|
|
registration_wait_timeout_seconds=0.5,
|
|
min_speech_ms=40,
|
|
trailing_silence_ms=40,
|
|
max_turn_ms=400,
|
|
asr_provider=_StubASRProvider(),
|
|
tts_provider=_ChunkedTTSProvider(),
|
|
load_registration_by_media_uuid=lambda value: None,
|
|
mark_media_connected=lambda session_id, value: None,
|
|
mark_media_ended=lambda session_id, reason: None,
|
|
touch_media_frame=lambda session_id: None,
|
|
set_state=lambda session_id, state, handoff_reason, metadata: None,
|
|
get_pending_greeting=lambda session_id: None,
|
|
mark_reply_delivered=lambda session_id, text, is_greeting: None,
|
|
plan_reply=lambda session_id, text, metadata, kind: None,
|
|
process_turn=lambda session_id, transcript_text, language, barge_in, metadata: None,
|
|
request_handoff=lambda session_id, customer_request_text, decision: None,
|
|
handle_media_error=lambda session_id, message, metadata: None,
|
|
)
|
|
runtime._tts_stream_prebuffer_ms = 40
|
|
|
|
writer = _FakeWriter()
|
|
|
|
async def _scenario() -> None:
|
|
actor = MediaActor(
|
|
registration=registration,
|
|
reader=asyncio.StreamReader(),
|
|
writer=writer, # type: ignore[arg-type]
|
|
vad=EnergyVAD(frame_ms=20, min_speech_ms=40, trailing_silence_ms=40, max_turn_ms=400),
|
|
frame_ms=20,
|
|
frame_bytes=320,
|
|
)
|
|
await runtime._speak_text(actor, "тест", is_greeting=False, reply_phase="main")
|
|
|
|
asyncio.run(_scenario())
|
|
|
|
assert len(writer.packets) > 1, "audio should be paced out frame by frame, not as one blob"
|
|
payload_bytes = 0
|
|
for packet in writer.packets:
|
|
packet_type, payload_length = struct.unpack("!BH", packet[:3])
|
|
assert packet_type == AUDIO_SOCKET_PACKET_PCM16
|
|
assert len(packet) == 3 + payload_length == 3 + 320, "every frame must be frame_bytes, zero-padded if short"
|
|
payload_bytes += payload_length
|
|
# 12 chunks * 160 bytes @16kHz downsample 2:1 -> 960 bytes @8kHz of real
|
|
# audio; frame padding on flush boundaries can only add silence, never drop it.
|
|
assert payload_bytes >= 960
|