feat: implement no-speech reprompt functionality with configurable parameters
deploy / deploy (push) Successful in 33s
deploy / deploy (push) Successful in 33s
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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())
|
||||
|
||||
Reference in New Issue
Block a user