Files
call-center/tests/test_ai_voice_media_runtime.py
T

1591 lines
63 KiB
Python

import asyncio
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_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,
)
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,
)
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,
)
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_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,
)
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,
)
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 _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