From c3992961658803249951f0e78bb56fcb6df07211 Mon Sep 17 00:00:00 2001 From: didar Date: Mon, 24 Aug 2026 14:17:00 +0500 Subject: [PATCH] feat: implement no-speech reprompt functionality with configurable parameters --- deployment/aimaq.env.production | 5 + .../ai_voice_runtime_service/media_runtime.py | 81 ++++++ tests/test_ai_voice_media_runtime.py | 270 ++++++++++++++++++ 3 files changed, 356 insertions(+) diff --git a/deployment/aimaq.env.production b/deployment/aimaq.env.production index 07fc800..24c9f91 100644 --- a/deployment/aimaq.env.production +++ b/deployment/aimaq.env.production @@ -59,6 +59,11 @@ AI_VOICE_VAD_RMS_THRESHOLD=800 AI_VOICE_VAD_MIN_SPEECH_MS=300 AI_VOICE_VAD_TRAILING_SILENCE_MS=500 +AI_VOICE_NO_SPEECH_REPROMPT_ENABLED=1 +AI_VOICE_NO_SPEECH_REPROMPT_SECONDS=7 +AI_VOICE_NO_SPEECH_REPROMPT_LOW_SIGNAL_TURNS=2 +AI_VOICE_NO_SPEECH_REPROMPT_MAX_ATTEMPTS=2 + AI_VOICE_ASR_PROVIDER=elevenlabs AI_VOICE_ASR_MODEL=gpt-4o-transcribe AI_VOICE_ASR_ELEVENLABS_API_KEY=sk_eb9bd36750e2a7fb378259b7144e978a044cfb412e6a2ac9 diff --git a/services/ai_voice_runtime_service/media_runtime.py b/services/ai_voice_runtime_service/media_runtime.py index e99db5d..f5f861a 100644 --- a/services/ai_voice_runtime_service/media_runtime.py +++ b/services/ai_voice_runtime_service/media_runtime.py @@ -122,6 +122,10 @@ class MediaActor: last_ack_variant: str | None = None current_reply_phase: str | None = None finalized_caller_turn_count: int = 0 + listening_since_monotonic: float = 0.0 + reprompt_attempts: int = 0 + consecutive_low_signal_turns: int = 0 + reprompt_watchdog_task: asyncio.Task | None = None class AudioSocketMediaRuntime: @@ -215,6 +219,20 @@ class AudioSocketMediaRuntime: str(os.getenv("AI_VOICE_V2_EARLY_PLAN_ENABLED", "1")).strip().lower() in {"1", "true", "yes", "on"} ) + self._no_speech_reprompt_enabled = ( + str(os.getenv("AI_VOICE_NO_SPEECH_REPROMPT_ENABLED", "1")).strip().lower() + in {"1", "true", "yes", "on"} + ) + self._no_speech_reprompt_seconds = max( + float(os.getenv("AI_VOICE_NO_SPEECH_REPROMPT_SECONDS", "7") or "7"), 2.0 + ) + self._no_speech_reprompt_low_signal_turns = max( + int(os.getenv("AI_VOICE_NO_SPEECH_REPROMPT_LOW_SIGNAL_TURNS", "2") or "2"), 1 + ) + self._no_speech_reprompt_max_attempts = max( + int(os.getenv("AI_VOICE_NO_SPEECH_REPROMPT_MAX_ATTEMPTS", "2") or "2"), 1 + ) + self._no_speech_reprompt_poll_seconds = 1.0 raw_early_plan_intents = str( os.getenv( "AI_VOICE_V2_EARLY_PLAN_INTENTS", @@ -340,6 +358,13 @@ class AudioSocketMediaRuntime: return "Сейчас уточню." return "Секунду." + @staticmethod + def _no_speech_reprompt_text(language: str | None) -> str: + normalized = str(language or "").strip().lower() + if normalized == "kz": + return "Алло, сізді тыңдап тұрмын. Немен көмектесе аламын?" + return "Алло, я вас слушаю. Подскажите, чем могу помочь?" + @staticmethod def _technical_issue_text(language: str | None) -> str: normalized = str(language or "").strip().lower() @@ -1316,6 +1341,7 @@ class AudioSocketMediaRuntime: await asyncio.to_thread(self._mark_media_connected, registration.voice_session_id, media_uuid) actor.keepalive_task = asyncio.create_task(self._keepalive_loop(actor)) actor.worker_task = asyncio.create_task(self._worker(actor)) + actor.reprompt_watchdog_task = asyncio.create_task(self._no_speech_reprompt_loop(actor)) while not actor.closed: packet_type, payload = await read_packet(reader, timeout=self._idle_timeout_seconds) @@ -1406,6 +1432,16 @@ class AudioSocketMediaRuntime: actor.registration.voice_session_id, ) await self._ensure_streaming_asr(actor) + if actor.asr_streaming_enabled and actor.asr_stream_id: + # The local VAD only flags speech_started after min_speech_ms of + # buffered pre-roll audio (see EnergyVAD.feed), and that pre-roll + # is folded into snapshot_utterance_pcm() together with the current + # frame. Without forwarding it here, every utterance's opening + # ~min_speech_ms would never reach the realtime ASR stream, clipping + # the first word(s) of each turn. + leadin_pcm = actor.vad.snapshot_utterance_pcm() + if len(leadin_pcm) > len(pcm_frame): + self._queue_streaming_asr_pcm(actor, leadin_pcm[: -len(pcm_frame)]) stream_input_active = actor.input_active or actor.barge_in_pending if actor.asr_streaming_enabled and actor.asr_stream_id and stream_input_active: self._queue_streaming_asr_pcm(actor, pcm_frame) @@ -1596,20 +1632,28 @@ class AudioSocketMediaRuntime: transcript_text[:160], ) if not transcript_text: + actor.consecutive_low_signal_turns += 1 await self._set_actor_state(actor, "listening") + if actor.consecutive_low_signal_turns >= self._no_speech_reprompt_low_signal_turns: + await self._maybe_play_no_speech_reprompt(actor, trigger="empty_transcript") return if self._is_low_signal_final_transcript(transcript_text): actor.finalized_utterance_generation = max(actor.finalized_utterance_generation, utterance_generation) + actor.consecutive_low_signal_turns += 1 logger.info( "audiosocket.low_signal_ignored session_id=%s transcript=%s", actor.registration.voice_session_id, transcript_text[:120], ) await self._set_actor_state(actor, "listening") + if actor.consecutive_low_signal_turns >= self._no_speech_reprompt_low_signal_turns: + await self._maybe_play_no_speech_reprompt(actor, trigger="low_signal_transcript") return actor.finalized_utterance_generation = max(actor.finalized_utterance_generation, utterance_generation) actor.finalized_caller_turn_count += 1 + actor.consecutive_low_signal_turns = 0 + actor.reprompt_attempts = 0 actor.partial_transcript = transcript_text if actor.registration.voice_v2_partial_asr else None actor.partial_intent = self._detect_early_intent(transcript_text) actor.stable_partial_intent = actor.partial_intent @@ -1925,6 +1969,37 @@ class AudioSocketMediaRuntime: ) await self._write_audio_packet(actor, silence_frame) + async def _maybe_play_no_speech_reprompt(self, actor: MediaActor, *, trigger: str) -> None: + if actor.closed or actor.state != "listening": + return + if actor.reprompt_attempts >= self._no_speech_reprompt_max_attempts: + return + actor.reprompt_attempts += 1 + actor.consecutive_low_signal_turns = 0 + logger.info( + "audiosocket.no_speech_reprompt session_id=%s attempt=%s trigger=%s", + actor.registration.voice_session_id, + actor.reprompt_attempts, + trigger, + ) + reprompt_text = self._no_speech_reprompt_text(actor.registration.language) + await self._speak_reply(actor, reprompt_text, is_greeting=False, reply_phase="no_speech_reprompt") + if not actor.closed: + await self._set_actor_state(actor, "listening") + + async def _no_speech_reprompt_loop(self, actor: MediaActor) -> None: + if not self._no_speech_reprompt_enabled: + return + while not actor.closed: + await asyncio.sleep(self._no_speech_reprompt_poll_seconds) + if actor.closed: + break + if actor.state != "listening" or actor.listening_since_monotonic <= 0: + continue + if time.monotonic() - actor.listening_since_monotonic < self._no_speech_reprompt_seconds: + continue + await self._maybe_play_no_speech_reprompt(actor, trigger="silence_timeout") + async def _set_actor_state( self, actor: MediaActor, @@ -1945,6 +2020,8 @@ class AudioSocketMediaRuntime: actor.state = state actor.input_active = state == "listening" actor.playback_active = state == "speaking" + if state == "listening": + actor.listening_since_monotonic = time.monotonic() await asyncio.to_thread( self._set_state, actor.registration.voice_session_id, @@ -1975,6 +2052,10 @@ class AudioSocketMediaRuntime: actor.keepalive_task.cancel() with contextlib.suppress(asyncio.CancelledError, Exception): await actor.keepalive_task + if actor.reprompt_watchdog_task is not None: + actor.reprompt_watchdog_task.cancel() + with contextlib.suppress(asyncio.CancelledError, Exception): + await actor.reprompt_watchdog_task if actor.handoff_task is not None: actor.handoff_task.cancel() with contextlib.suppress(asyncio.CancelledError, Exception): diff --git a/tests/test_ai_voice_media_runtime.py b/tests/test_ai_voice_media_runtime.py index 112d343..3307987 100644 --- a/tests/test_ai_voice_media_runtime.py +++ b/tests/test_ai_voice_media_runtime.py @@ -1,4 +1,5 @@ import asyncio +import contextlib import threading import time import uuid @@ -2629,3 +2630,272 @@ def _legacy_test_media_runtime_voice_v2_falls_back_when_streaming_sidecar_is_una 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 + runtime._no_speech_reprompt_poll_seconds = 0.02 + + 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, + ) + await runtime._set_actor_state(actor, "listening") + watchdog_task = asyncio.create_task(runtime._no_speech_reprompt_loop(actor)) + for _ in range(50): + if spoken: + break + await asyncio.sleep(0.02) + actor.closed = True + watchdog_task.cancel() + with contextlib.suppress(asyncio.CancelledError, Exception): + await watchdog_task + assert spoken, "watchdog should reprompt after prolonged silence in listening state" + + asyncio.run(_scenario())