from __future__ import annotations import asyncio import importlib import json import os import sys import time import types import unittest import uuid from pathlib import Path from tempfile import TemporaryDirectory from typing import Any from unittest import mock from types import SimpleNamespace def _install_fake_livekit() -> None: try: importlib.import_module("livekit.agents") return except ImportError: pass livekit_module = types.ModuleType("livekit") agents_module = types.ModuleType("livekit.agents") class UserInputTranscribedEvent: def __init__(self, transcript: str = "", is_final: bool = True) -> None: self.transcript = transcript self.is_final = is_final class _AudioInputOptions: def __init__(self, **kwargs) -> None: self.kwargs = kwargs class _AudioOutputOptions: def __init__(self, **kwargs) -> None: self.kwargs = kwargs class _RoomOptions: def __init__(self, **kwargs) -> None: self.kwargs = kwargs agents_module.UserInputTranscribedEvent = UserInputTranscribedEvent agents_module.room_io = SimpleNamespace( AudioInputOptions=_AudioInputOptions, AudioOutputOptions=_AudioOutputOptions, RoomOptions=_RoomOptions, ) livekit_module.agents = agents_module sys.modules["livekit"] = livekit_module sys.modules["livekit.agents"] = agents_module _install_fake_livekit() from app.livekit.policies.finalization_policy import FinalizationPolicy from app.livekit.policies.idle_policy import IdlePolicy, IdlePolicyConfig from app.livekit.policies.interrupt_policy import InterruptPolicy from app.livekit.adapters.speech_service import SpeechService from app.livekit.adapters.agent_backend import BackendReply from app.livekit.runtime.call_runtime import ( CallRuntime, InflightBackendWaitNoticeResult, InflightBackendWaitTimedOut, ) from app.livekit.runtime.command_executor import RuntimeCommandExecutor from app.livekit.runtime.commands import ( EndServiceOnce, StartSpeech, ExportSession, ExtractSpokenText, InjectIdleNudge, InterruptSpeech, NotifyBridgeDone, NotifyBridgeStop, RunPipelineInput, SetPendingInterrupt, SetPipelineInterruption, SetPipelineProcessingInterruption, StartSession, ) from app.livekit.runtime.scheduler import TimerScheduler from app.livekit.runtime.state import RuntimeConfig from app.livekit.vad_dynamic_threshold import ( DynamicVADThresholdConfig, DynamicVADThresholdController, ) from app.providers.stt_internal_livekit import stt_text_with_single_word_threshold from app.utils.logging import StructuredLogContext from app.utils.turn_ids import ( peek_started_turn_message_id, register_started_turn_message_id, register_transcribed_turn, reset_turn_message_sequence, ) WAIT_LONG_AUDIO_PATH = str( Path(__file__).resolve().parents[2] / "src" / "app" / "livekit" / "assets" / "comfort" / "long" / "01.wav" ) DEFAULT_WAIT_TEXT = "Um momento, ainda estou consultando para te ajudar." def _expected_wait_text_for_audio(path: str | Path, fallback: str = DEFAULT_WAIT_TEXT) -> str: text_path = Path(path).with_suffix(".txt") if not text_path.is_file(): return fallback for encoding in ("utf-8-sig", "utf-8", "latin-1"): try: text = " ".join(text_path.read_text(encoding=encoding).split()) except UnicodeDecodeError: continue except OSError: return fallback return text or fallback return fallback class _NullLogger: def info(self, *args, **kwargs) -> None: return None def debug(self, *args, **kwargs) -> None: return None def warning(self, *args, **kwargs) -> None: return None def exception(self, *args, **kwargs) -> None: return None class _AsyncTestCase(unittest.IsolatedAsyncioTestCase): async def asyncSetUp(self) -> None: asyncio.get_running_loop().set_debug(False) class _FakeAgent: def __init__(self) -> None: self._ready = asyncio.Event() self._ready.set() self._run_lock = asyncio.Lock() self.pipeline = object() self.pending_interrupts: list[tuple[str, bool, str]] = [] self._consumed_interrupt = (None, "", False) self.end_calls = 0 async def set_pending_interrupt( self, *, listened_text: str = "", skipped: bool = False, speech_id: str = "", ) -> None: self.pending_interrupts.append((listened_text, skipped, speech_id)) async def consume_pending_interrupt(self): return self._consumed_interrupt async def end_service_once(self): self.end_calls += 1 return BackendReply(stage="DONE", done=True, export_payload=[{"status": "ok"}]) class _FakePushPipeline: def __init__(self, replies) -> None: self._replies = asyncio.Queue() for reply in replies: self._replies.put_nowait(reply) def supports_server_push(self) -> bool: return True async def wait_for_server_push(self): if self._replies.empty(): return None return await self._replies.get() class _FakeInflightPushPipeline: def __init__(self) -> None: self._replies = asyncio.Queue() self._closed = asyncio.Event() def supports_inflight_backend_push(self) -> bool: return True def supports_server_push(self) -> bool: return False async def push(self, reply) -> None: await self._replies.put(reply) def close(self) -> None: self._closed.set() async def wait_for_server_push(self): if not self._replies.empty(): return await self._replies.get() if self._closed.is_set(): return None get_task = asyncio.create_task(self._replies.get()) close_task = asyncio.create_task(self._closed.wait()) done, _pending = await asyncio.wait( {get_task, close_task}, return_when=asyncio.FIRST_COMPLETED, ) if get_task in done: close_task.cancel() await asyncio.gather(close_task, return_exceptions=True) return get_task.result() get_task.cancel() await asyncio.gather(get_task, return_exceptions=True) return None class _FakeAgentInput: def __init__(self) -> None: self.audio_enabled = True self.calls = [] def set_audio_enabled(self, enabled: bool) -> None: self.audio_enabled = bool(enabled) self.calls.append(bool(enabled)) class _FakeSession: def __init__(self) -> None: self.handlers = {} self.start_calls = [] self.input = _FakeAgentInput() def on(self, event_name: str): def _decorator(fn): self.handlers[event_name] = fn return fn return _decorator async def start(self, **kwargs) -> None: self.start_calls.append(kwargs) class _FakeRoom: def __init__(self) -> None: self.name = "room-test" self.remote_participants = {} self.handlers = {} def on(self, event_name: str): def _decorator(fn): self.handlers[event_name] = fn return fn return _decorator class _FakeTimeline: def __init__(self) -> None: self.events = [] def emit(self, event: str, **fields) -> None: self.events.append((event, fields)) class _FakeContext: def __init__(self) -> None: self.room = _FakeRoom() self.shutdown_callbacks = [] def add_shutdown_callback(self, callback) -> None: self.shutdown_callbacks.append(callback) class _FakeCommandExecutor: def __init__(self) -> None: self.commands = [] self.pipeline_result = BackendReply(stage="PRESENTATION", text="resposta") self.end_reply = BackendReply(stage="DONE", done=True, export_payload=[{"status": "ok"}]) self.spoken_text = "trecho falado" self.speech_handle = SimpleNamespace( id="speech-test", interrupted=False, ) self.speech_handles = None self.run_pipeline_side_effect = None self.on_start_session: Any = None self.wait_for_playout_delay_s = 0.0 self.wait_for_playout_delays_s = [] self.wait_for_playout_exceptions = [] self.wait_for_playout_callbacks = [] async def execute(self, command): self.commands.append(command) if isinstance(command, RunPipelineInput): if self.run_pipeline_side_effect is not None: await self.run_pipeline_side_effect() return self.pipeline_result if isinstance(command, EndServiceOnce): return self.end_reply if isinstance(command, ExportSession): return None if isinstance(command, NotifyBridgeDone): return None if isinstance(command, NotifyBridgeStop): return None if isinstance(command, SetPendingInterrupt): return None if command.__class__.__name__ == "StartSpeech": if self.speech_handles: return self.speech_handles.pop(0) return self.speech_handle if command.__class__.__name__ == "WaitForSpeechPlayout": if self.wait_for_playout_callbacks: callback = self.wait_for_playout_callbacks.pop(0) if callback is not None: result = callback(command) if asyncio.iscoroutine(result): await result if self.wait_for_playout_exceptions: exc = self.wait_for_playout_exceptions.pop(0) if exc is not None: raise exc delay_s = ( self.wait_for_playout_delays_s.pop(0) if self.wait_for_playout_delays_s else self.wait_for_playout_delay_s ) if delay_s > 0: await asyncio.sleep(delay_s) return None if isinstance(command, InterruptSpeech): setattr(command.handle, "interrupted", True) return None if command.__class__.__name__ == "SetPipelineInterruption": return None if command.__class__.__name__ == "InjectIdleNudge": return None if command.__class__.__name__ == "StartSession": if self.on_start_session is not None: self.on_start_session(command) return None raise AssertionError(f"Unsupported command in fake executor: {command!r}") def execute_now(self, command): self.commands.append(command) if isinstance(command, ExtractSpokenText): return self.spoken_text raise AssertionError(f"Unsupported synchronous command: {command!r}") class _FakeVADThresholdController: def __init__(self) -> None: self.calls = [] self.active = False def activate_agent_wait_timeout_retry(self, *, attempt: int, reason: str) -> bool: if self.active: return False self.calls.append(("activate", attempt, reason)) self.active = True return True def restore(self, *, reason: str) -> bool: if not self.active: return False self.calls.append(("restore", reason)) self.active = False return True class _FakeVAD: def __init__(self) -> None: self.update_options_calls = [] def update_options(self, **kwargs) -> None: self.update_options_calls.append(kwargs) class _AwaitableSpeechHandle: def __init__(self) -> None: self.awaited = False async def wait_for_playout(self) -> None: return None def interrupt(self, *, force: bool = False) -> None: return None def __await__(self): async def _wait(): self.awaited = True return self return _wait().__await__() class SpeechServiceTests(_AsyncTestCase): async def test_start_returns_livekit_speech_handle_without_waiting_for_playout(self) -> None: handle = _AwaitableSpeechHandle() session = SimpleNamespace( say=lambda *args, **kwargs: handle, ) service = SpeechService(session) result = await service.start( "texto", allow_interruptions=True, add_to_chat_ctx=True, ) self.assertIs(result, handle) self.assertFalse(handle.awaited) async def test_start_passes_provided_audio_to_livekit_session(self) -> None: handle = _AwaitableSpeechHandle() calls = [] audio = object() session = SimpleNamespace( say=lambda *args, **kwargs: calls.append((args, kwargs)) or handle, ) service = SpeechService(session) result = await service.start( "texto", allow_interruptions=True, add_to_chat_ctx=False, audio=audio, ) self.assertIs(result, handle) self.assertIs(calls[0][1]["audio"], audio) self.assertFalse(calls[0][1]["add_to_chat_ctx"]) class DynamicVADThresholdControllerTests(unittest.TestCase): def test_activate_and_restore_updates_vad_options(self) -> None: vad = _FakeVAD() controller = DynamicVADThresholdController( vad, config=DynamicVADThresholdConfig( enabled=True, baseline_activation_threshold=0.30, baseline_deactivation_threshold=0.15, retry_activation_threshold=0.20, retry_deactivation_threshold=0.05, ), ) self.assertEqual( controller.current_thresholds(), { "mode": "baseline", "activation_threshold": 0.30, "deactivation_threshold": 0.15, }, ) self.assertTrue( controller.activate_agent_wait_timeout_retry( attempt=1, reason="agent_wait_timeout_retry", ) ) self.assertEqual( controller.current_thresholds(), { "mode": "agent_wait_timeout_retry", "activation_threshold": 0.20, "deactivation_threshold": 0.05, }, ) self.assertTrue(controller.restore(reason="user_final")) self.assertEqual( controller.current_thresholds(), { "mode": "baseline", "activation_threshold": 0.30, "deactivation_threshold": 0.15, }, ) self.assertEqual( vad.update_options_calls, [ { "activation_threshold": 0.20, "deactivation_threshold": 0.05, }, { "activation_threshold": 0.30, "deactivation_threshold": 0.15, }, ], ) class InterruptPolicyTests(unittest.TestCase): def test_allow_stage_and_backchannel(self) -> None: policy = InterruptPolicy() self.assertTrue(policy.allow_stage("ARGUMENTATION")) self.assertFalse(policy.allow_stage("FORMALIZATION")) self.assertFalse(policy.allow_stage("IDLE_NUDGE")) self.assertTrue(policy.allow_stage("AGENT_WAIT_TIMEOUT_RETRY")) self.assertTrue(policy.allow_stage("unknown-stage")) self.assertTrue(policy.should_ignore_backchannel(speaking=True, text="uhum")) self.assertFalse(policy.should_ignore_backchannel(speaking=False, text="uhum")) class IdlePolicyTests(unittest.TestCase): def test_idle_policy_decisions(self) -> None: policy = IdlePolicy( IdlePolicyConfig( enabled=True, delay_s=10.0, join_delay_s=15.0, close_delay_s=20.0, max_tries=3, end_reason="no_user_response", ) ) self.assertTrue(policy.should_arm_nudge(delay_s=10.0, nudge_text="alo")) self.assertFalse(policy.should_arm_nudge(delay_s=0.0, nudge_text="alo")) self.assertTrue( policy.should_fire_nudge( token_is_current=True, seq_matches=True, speaking=False, gap_active=False, finalized=False, current_stage="PRESENTATION", ) ) self.assertFalse( policy.should_fire_nudge( token_is_current=True, seq_matches=True, speaking=False, gap_active=False, finalized=False, current_stage="DONE", ) ) self.assertTrue(policy.should_arm_close_after_nudge(stage="IDLE_NUDGE", nudge_count=3)) def test_idle_policy_is_disabled_by_default(self) -> None: policy = IdlePolicy( IdlePolicyConfig( enabled=False, delay_s=10.0, join_delay_s=15.0, close_delay_s=20.0, max_tries=3, end_reason="no_user_response", ) ) self.assertFalse(policy.should_arm_nudge(delay_s=10.0, nudge_text="alo")) self.assertFalse(policy.should_arm_close_before_fire(nudge_count=3)) self.assertFalse(policy.should_arm_close_after_nudge(stage="IDLE_NUDGE", nudge_count=3)) class FinalizationPolicyTests(unittest.TestCase): def test_finalization_decisions(self) -> None: policy = FinalizationPolicy(call_end_grace_s=2.0) self.assertTrue(policy.should_skip(finalized=True)) self.assertFalse(policy.should_skip(finalized=False)) self.assertTrue(policy.should_finalize_room_empty(remote_participants=0, finalized=False)) self.assertFalse(policy.should_finalize_room_empty(remote_participants=1, finalized=False)) self.assertTrue(policy.is_done_stage("DONE")) self.assertFalse(policy.is_done_stage("PRESENTATION")) class TimerSchedulerTests(_AsyncTestCase): async def test_arm_updates_token_and_cancel_stops_task(self) -> None: created_tasks = [] logger = _NullLogger() def _create_task_logged(coro, *, name: str): task = asyncio.create_task(coro, name=name) created_tasks.append(task) return task scheduler = TimerScheduler(create_task_logged=_create_task_logged, logger=logger) gate = asyncio.Event() async def _wait_forever(): await gate.wait() token1 = scheduler.arm("idle", task_name="idle_1", coro=_wait_forever()) token2 = scheduler.arm("idle", task_name="idle_2", coro=_wait_forever()) self.assertFalse(scheduler.is_current("idle", token1)) self.assertTrue(scheduler.is_current("idle", token2)) scheduler.cancel("idle", reason="test", log_name="IDLE_TIMER_CANCEL") await asyncio.sleep(0) self.assertTrue(created_tasks[-1].cancelled()) gate.set() await asyncio.gather(*created_tasks, return_exceptions=True) class RuntimeCommandExecutorTests(_AsyncTestCase): async def test_routes_commands_to_underlying_services(self) -> None: agent = _FakeAgent() bridge_gateway = SimpleNamespace(notify_stage_done=self._make_async_noop()) bridge_gateway.notify_stop = self._make_async_noop() export_service = SimpleNamespace(export_session=self._make_async_noop()) session = _FakeSession() speech_service = SimpleNamespace( start=self._make_async_return("speech-handle"), wait_for_playout=self._make_async_noop(), extract_spoken_text=lambda handle: "spoken", ) agent.pipeline = SimpleNamespace( inject_idle_nudge=self._make_async_noop(), set_interruption=self._make_async_noop(), run=self._make_async_return(BackendReply(stage="DONE", text="encerrando", done=True)), ) executor = RuntimeCommandExecutor( agent=agent, bridge_gateway=bridge_gateway, export_service=export_service, session=session, speech_service=speech_service, ) result = await executor.execute(RunPipelineInput("payload")) self.assertEqual(result, BackendReply(stage="DONE", text="encerrando", done=True)) self.assertEqual(executor.execute_now(ExtractSpokenText(object())), "spoken") @staticmethod def _make_async_noop(): async def _noop(*args, **kwargs): return None return _noop @staticmethod def _make_async_return(value): async def _return(*args, **kwargs): return value return _return class CallRuntimeTests(_AsyncTestCase): def setUp(self) -> None: self.created_tasks = [] def _create_task_logged(self, coro, *, name: str): task = asyncio.create_task(coro, name=name) self.created_tasks.append(task) return task async def _drain_tasks(self) -> None: while True: pending = [task for task in self.created_tasks if not task.done()] if not pending: return await asyncio.gather(*pending) async def _cancel_pending_tasks(self) -> None: pending = [task for task in self.created_tasks if not task.done()] for task in pending: task.cancel() if pending: await asyncio.gather(*pending, return_exceptions=True) def test_inflight_backend_wait_message_id_uses_uuid_when_original_is_missing(self) -> None: generated = uuid.UUID("12345678-1234-4234-9234-123456789abc") with mock.patch("app.livekit.runtime.call_runtime.uuid.uuid4", return_value=generated): message_id = CallRuntime._inflight_backend_wait_message_id(attempt=1) self.assertEqual( message_id, "12345678-1234-4234-9234-123456789abc_conforto_12345678123442349234123456789abc", ) def test_feedback_message_id_keeps_the_agent_speech_id(self) -> None: message_id = CallRuntime._feedback_message_id( message_id="turno-1", speech_id="speech-feedback-1", ) self.assertEqual(message_id, "turno-1_feedback_speech-feedback-1") self.assertEqual(CallRuntime._base_message_id(message_id), "turno-1") def test_auxiliary_message_ids_are_distinct_and_keep_the_base(self) -> None: generated = uuid.UUID("12345678-1234-4234-9234-123456789abc") with mock.patch("app.livekit.runtime.call_runtime.uuid.uuid4", return_value=generated): interruption = CallRuntime._interruption_comfort_message_id(message_id="turno-1") tts_error = CallRuntime._tts_error_message_id(message_id="turno-1") self.assertEqual(CallRuntime._idle_nudge_message_id(message_id="turno-1", attempt=2), "turno-1_idle_nudge_2") self.assertEqual(interruption, "turno-1_interruption_confort_12345678123442349234123456789abc") self.assertEqual(tts_error, "turno-1_tts_error_12345678123442349234123456789abc") self.assertEqual(CallRuntime._base_message_id(interruption), "turno-1") self.assertEqual(CallRuntime._base_message_id(tts_error), "turno-1") self.assertTrue( CallRuntime._transfer_message_id(message_id="turno-1").startswith( "turno-1_transferencia_" ) ) self.assertTrue( CallRuntime._resource_error_message_id( message_id="turno-1", terminal=True, tipo_evento="envio msg", ).startswith("turno-1_erro_terminal_envio_") ) async def test_backend_feedback_uses_a_suffixed_message_id(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} with mock.patch.object(runtime, "say_stage", new_callable=mock.AsyncMock) as say_stage: await runtime._speak_backend_reply( BackendReply( stage="PRESENTATION", text="Ainda estou consultando.", metadata={ "message_id": "turno-feedback", "speech_id": "speech-feedback-1", "agent_message_type": "feedback", }, ), source="backend_push", add_to_chat_ctx=True, ) self.assertEqual( say_stage.await_args.kwargs["message_id"], "turno-feedback_feedback_speech-feedback-1", ) self.assertEqual(say_stage.await_args.kwargs["speech_id"], "speech-feedback-1") def _make_runtime( self, *, agent_starts_conversation: bool = False, push_replies=None, runtime_config_overrides=None, idle_policy_overrides=None, vad_threshold_controller=None, ): agent = _FakeAgent() if push_replies is not None: agent.pipeline = _FakePushPipeline(push_replies) session = _FakeSession() ctx = _FakeContext() executor = _FakeCommandExecutor() config_kwargs = { "call_end_grace_s": 0.0, "final_grace_s": 0.0, "idle_nudge_delay_s": 0.0, "idle_nudge_join_delay_s": 0.0, "idle_nudge_close_delay_s": 0.0, "idle_nudge_max_tries": 0, "idle_nudge_end_reason": "no_user_response", } if runtime_config_overrides: config_kwargs.update(runtime_config_overrides) idle_policy_kwargs = { "enabled": False, "delay_s": 0.0, "join_delay_s": 0.0, "close_delay_s": 0.0, "max_tries": 0, "end_reason": "no_user_response", } if idle_policy_overrides: idle_policy_kwargs.update(idle_policy_overrides) runtime = CallRuntime( ctx=ctx, session=session, agent=agent, command_executor=executor, call_logger=_NullLogger(), protocol="PRT-1", session_id="S-1", bridge_identity="bridge-1", agent_starts_conversation=agent_starts_conversation, nudge_text="alo", call_t0=0.0, extract_text_from_transcript=lambda transcript: transcript, interrupt_policy=InterruptPolicy(), idle_policy=IdlePolicy(IdlePolicyConfig(**idle_policy_kwargs)), finalization_policy=FinalizationPolicy(call_end_grace_s=0.0), create_task_logged=self._create_task_logged, config=RuntimeConfig(**config_kwargs), logger=_NullLogger(), vad_threshold_controller=vad_threshold_controller, ) return runtime, agent, executor def test_inflight_backend_wait_uses_short_audio_then_long_audio(self) -> None: with TemporaryDirectory() as tmp_dir: base_dir = Path(tmp_dir) short_dir = base_dir / "short" long_dir = base_dir / "long" short_dir.mkdir() long_dir.mkdir() short_audio = short_dir / "curto.wav" long_audio = long_dir / "longo.wav" short_audio.write_bytes(b"short") long_audio.write_bytes(b"long") runtime, _agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_short_audio_dir": str(short_dir), "inflight_backend_wait_long_audio_dir": str(long_dir), } ) self.assertEqual( runtime._inflight_backend_wait_audio_path(attempt=1), short_audio, ) self.assertEqual( runtime._inflight_backend_wait_audio_path(attempt=2), long_audio, ) def test_inflight_backend_wait_prewarms_audio_duration_cache(self) -> None: with TemporaryDirectory() as tmp_dir: short_dir = Path(tmp_dir) / "short" short_dir.mkdir() short_audio = short_dir / "curto.wav" short_audio.write_bytes(b"fake-wav") with mock.patch( "app.livekit.runtime.call_runtime.wav_duration_ms", return_value=321, ) as duration_ms: runtime, _agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_short_audio_dir": str(short_dir), } ) duration_ms.assert_called_once_with(str(short_audio)) with mock.patch( "app.livekit.runtime.call_runtime.wav_duration_ms", side_effect=AssertionError("duration should come from cache"), ): self.assertEqual( runtime._inflight_backend_wait_audio_duration_ms(short_audio), 321, ) def test_inflight_backend_wait_keeps_legacy_long_audio_path_fallback(self) -> None: runtime, _agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) self.assertEqual( runtime._inflight_backend_wait_audio_path(), Path(WAIT_LONG_AUDIO_PATH), ) async def test_idle_nudge_does_not_activate_retry_vad_threshold(self) -> None: controller = _FakeVADThresholdController() runtime, _agent, executor = self._make_runtime( idle_policy_overrides={ "enabled": True, "delay_s": 10.0, "join_delay_s": 10.0, "close_delay_s": 10.0, "max_tries": 2, "end_reason": "no_user_response", }, vad_threshold_controller=controller, ) try: runtime.arm_idle_timer(reason="test", delay_s=0.001) for _ in range(50): if any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands): break await asyncio.sleep(0.001) self.assertNotIn(("activate", 1, "idle_nudge_fire"), controller.calls) start_speech = next( cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ) self.assertEqual(start_speech.text, "alo") self.assertFalse(start_speech.allow_interruptions) runtime.on_user_input_transcribed(SimpleNamespace(transcript="sim", is_final=True)) self.assertNotIn(("restore", "user_final"), controller.calls) finally: await self._cancel_pending_tasks() async def test_idle_nudge_uses_a_suffixed_message_id(self) -> None: runtime, _agent, _executor = self._make_runtime( idle_policy_overrides={ "enabled": True, "delay_s": 10.0, "join_delay_s": 10.0, "close_delay_s": 10.0, "max_tries": 2, "end_reason": "no_user_response", } ) runtime._last_agent_message_id = "turno-idle" try: with mock.patch.object(runtime, "say_stage", new_callable=mock.AsyncMock) as say_stage: runtime.arm_idle_timer(reason="test", delay_s=0.001) for _ in range(50): if say_stage.await_count: break await asyncio.sleep(0.001) say_stage.assert_awaited_once() self.assertEqual( say_stage.await_args.kwargs["message_id"], "turno-idle_idle_nudge_1", ) finally: await self._cancel_pending_tasks() async def test_on_user_input_transcribed_stashes_when_stage_is_not_interruptible(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "FORMALIZATION" runtime.speaking.set() event = SimpleNamespace(transcript="quero continuar", is_final=True) with mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log: runtime.on_user_input_transcribed(event) await self._drain_tasks() self.assertIsNotNone(runtime.state.pending_user_final) self.assertEqual(runtime.state.pending_user_final.text, "quero continuar") self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( call.args[1] == "interrupt_ignored" and call.kwargs["reason"] == "stage_not_interruptible" for call in flow_log.call_args_list ) ) async def test_on_user_input_transcribed_runs_when_argumentation_is_speaking(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "ARGUMENTATION" runtime.state.current_speech.allow_interruptions = True runtime.speaking.set() event = SimpleNamespace(transcript="quero continuar", is_final=True) runtime.on_user_input_transcribed(event) await self._drain_tasks() self.assertIsNone(runtime.state.pending_user_final) self.assertTrue(any(isinstance(cmd, SetPendingInterrupt) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_short_vad_utterance_with_text_is_discarded_while_r1_is_processing(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True register_transcribed_turn( runtime._structured_log_context, message_id="short-utterance", transcription="sim", text="sim", ) runtime.note_vad_speech_end(999) runtime.on_user_input_transcribed(SimpleNamespace(transcript="sim", is_final=True)) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 1) self.assertEqual(runtime.state.deferred_interruption.long_turns, []) self.assertFalse(runtime.state.deferred_interruption.special_comfort_sent) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_first_long_vad_plays_special_before_empty_stt(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True with mock.patch.object(runtime, "_play_deferred_interruption_comfort") as comfort: runtime.note_vad_speech_end(1000) comfort.assert_called_once_with() runtime.on_user_input_transcribed(SimpleNamespace(transcript="", is_final=True)) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 1) self.assertEqual(runtime.state.deferred_interruption.long_turns, []) self.assertTrue(runtime.state.deferred_interruption.special_comfort_sent) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_special_processing_comfort_starts_protected(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.deferred_interruption.special_comfort_sent = True runtime._play_deferred_interruption_comfort() await self._drain_tasks() speech = next(cmd for cmd in executor.commands if isinstance(cmd, StartSpeech)) self.assertEqual(speech.text, "Um instante") self.assertFalse(speech.allow_interruptions) async def test_short_backend_wait_comfort_starts_protected(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} with mock.patch.object( runtime, "_select_inflight_backend_wait_audio", return_value=(Path(WAIT_LONG_AUDIO_PATH), "short"), ): played = await runtime._play_inflight_backend_wait_notice_audio( user_seq=None, attempt=1, ) self.assertTrue(played) speech = next(cmd for cmd in executor.commands if isinstance(cmd, StartSpeech)) self.assertFalse(speech.allow_interruptions) async def test_short_comfort_is_not_rearmed_by_its_originating_user_final(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.speaking.set() runtime.state.current_speech.stage = "AGENT_BACKEND_WAIT" runtime.state.current_speech.handle = executor.speech_handle runtime.state.current_speech.allow_interruptions = False runtime.state.current_speech.rearm_interruptions_on_next_user_speech = True executor.speech_handle.allow_interruptions = False runtime.on_user_input_transcribed( SimpleNamespace(transcript="fala original", is_final=True) ) await self._drain_tasks() self.assertFalse(runtime.state.current_speech.allow_interruptions) self.assertFalse(executor.speech_handle.allow_interruptions) self.assertTrue( runtime.state.current_speech.rearm_interruptions_on_next_user_speech ) async def test_short_comfort_is_rearmed_on_new_user_speech(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.speaking.set() runtime.state.current_speech.stage = "INTERRUPTION_COMFORT" runtime.state.current_speech.handle = executor.speech_handle runtime.state.current_speech.allow_interruptions = False runtime.state.current_speech.rearm_interruptions_on_next_user_speech = True executor.speech_handle.allow_interruptions = False runtime._rearm_current_speech_interruptions_on_new_user_speech() self.assertTrue(runtime.state.current_speech.allow_interruptions) self.assertFalse( runtime.state.current_speech.rearm_interruptions_on_next_user_speech ) self.assertTrue(executor.speech_handle.allow_interruptions) async def test_feedback_compat_mode_drops_processing_speech_without_replay(self) -> None: runtime, _agent, executor = self._make_runtime( runtime_config_overrides={"deferred_interruption_enabled": False} ) runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True register_transcribed_turn( runtime._structured_log_context, message_id="processing-drop", transcription="quero interromper", text="quero interromper", ) with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort") as special, mock.patch.object(runtime, "_play_deferred_short_comfort") as short, mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log, ): runtime.note_vad_speech_end(2500) runtime.on_user_input_transcribed( SimpleNamespace(transcript="quero interromper", is_final=True) ) await self._drain_tasks() special.assert_not_called() short.assert_not_called() self.assertEqual(runtime.state.user_final_seq, 1) self.assertEqual(runtime.state.deferred_interruption.long_turns, []) self.assertEqual(len(runtime.state.deferred_interruption.pending_vad_utterances), 0) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( call.args[1] == "backend_processing_vad_discarded" and call.kwargs["mode"] == "feedback_compat" for call in flow_log.call_args_list ) ) self.assertTrue( any( call.args[1] == "user_input_dropped" and call.kwargs["reason"] == "backend_processing_feedback_compat" and call.kwargs["agent_message_type"] == "feedback" for call in flow_log.call_args_list ) ) async def test_vad_end_outside_backend_window_is_not_deferred(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.state.deferred_interruption.backend_in_flight = False runtime.speaking.set() with mock.patch.object(runtime, "_play_deferred_interruption_comfort") as comfort: runtime.note_vad_speech_end(4000) comfort.assert_not_called() self.assertEqual(len(runtime.state.deferred_interruption.pending_vad_utterances), 0) self.assertFalse(runtime.state.deferred_interruption.special_comfort_sent) async def test_barge_in_during_agent_playout_still_runs_a_new_pipeline(self) -> None: runtime, agent, executor = self._make_runtime() runtime.state.current_stage = "ARGUMENTATION" runtime.state.current_speech.allow_interruptions = True runtime.speaking.set() # O playout segura o run_lock, mas a janela do ciclo diferido ja fechou: # a fala precisa voltar a ser barge-in comum. await agent._run_lock.acquire() try: runtime.note_vad_speech_end(4000) runtime.on_user_input_transcribed( SimpleNamespace(transcript="muda de assunto", is_final=True) ) await asyncio.sleep(0) finally: agent._run_lock.release() await self._drain_tasks() self.assertEqual(runtime.state.deferred_interruption.long_turns, []) self.assertTrue(any(isinstance(cmd, SetPendingInterrupt) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_long_vad_transcripts_are_concatenated_once_and_r2_blocks_replay(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort") as comfort, mock.patch.object(runtime, "_play_deferred_short_comfort") as short_comfort, mock.patch.object(runtime, "run_pipeline", new_callable=mock.AsyncMock) as run_pipeline, ): runtime.note_vad_speech_end(1200) runtime.note_vad_speech_end(300) runtime.note_vad_speech_end(1400) register_transcribed_turn( runtime._structured_log_context, message_id="long-1", transcription="primeira fala longa", text="primeira fala longa", ) register_transcribed_turn( runtime._structured_log_context, message_id="short-1", transcription="fala curta", text="fala curta", ) register_transcribed_turn( runtime._structured_log_context, message_id="long-2", transcription="segunda fala longa", text="segunda fala longa", ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="primeira fala longa", is_final=True) ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="fala curta", is_final=True) ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="segunda fala longa", is_final=True) ) comfort.assert_called_once_with() short_comfort.assert_called_once_with(source="repeated_interruption") self.assertEqual( [turn.text for turn in runtime.state.deferred_interruption.long_turns], ["primeira fala longa", "segunda fala longa"], ) dispatched = await runtime._dispatch_deferred_interruption( BackendReply(stage="PRESENTATION", text="resposta R1") ) await self._drain_tasks() self.assertTrue(dispatched) self.assertTrue(runtime.state.deferred_interruption.replay_in_flight) run_pipeline.assert_awaited_once_with( "primeira fala longa. segunda fala longa", "primeira fala longa. segunda fala longa", is_deferred_replay=True, user_seq=3, message_id="long-2", inflight_initial_notice_consumed=True, ) runtime.note_vad_speech_end(1500) runtime.on_user_input_transcribed( SimpleNamespace(transcript="nao deve gerar R3", is_final=True) ) self.assertEqual(runtime.state.user_final_seq, 3) self.assertEqual(runtime.state.deferred_interruption.long_turns, []) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_vad_ended_before_r1_is_promoted_when_stt_finishes_during_r1(self) -> None: runtime, _agent, executor = self._make_runtime() backend_started = asyncio.Event() release_backend = asyncio.Event() runtime.state.user_final_seq = 0 async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() executor.run_pipeline_side_effect = _wait_for_backend with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort"), mock.patch.object( runtime, "_dispatch_deferred_interruption", new_callable=mock.AsyncMock, return_value=True, ) as dispatch, ): # O segundo fim de fala acontece antes de R1, mas seu STT chega # somente enquanto R1 está em voo. runtime.note_vad_speech_end(1500) runtime.note_vad_speech_end(1500) register_transcribed_turn( runtime._structured_log_context, message_id="r1-turn", transcription="primeiro turno", text="primeiro turno", ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="primeiro turno", is_final=True) ) await backend_started.wait() register_transcribed_turn( runtime._structured_log_context, message_id="late-turn", transcription="fala que terminou antes de R1", text="fala que terminou antes de R1", ) runtime.on_user_input_transcribed( SimpleNamespace( transcript="fala que terminou antes de R1", is_final=True ) ) release_backend.set() # O task criado pelo callback pode precisar de um ciclo extra para # concluir a espera e tentar o dispatch. await self._drain_tasks() dispatch.assert_awaited_once() self.assertEqual(len(runtime.state.deferred_interruption.pre_backend_vad_utterances), 0) async def test_r1_waits_for_pending_long_stt_before_dispatching(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 backend_started = asyncio.Event() release_backend = asyncio.Event() async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() executor.run_pipeline_side_effect = _wait_for_backend with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort"), mock.patch.object(runtime, "_dispatch_deferred_interruption", new_callable=mock.AsyncMock, return_value=True) as dispatch, ): task = asyncio.create_task(runtime.run_pipeline("original", "original", user_seq=1)) await backend_started.wait() runtime.note_vad_speech_end(1200) release_backend.set() await asyncio.sleep(0) self.assertFalse(task.done()) self.assertFalse(any(isinstance(cmd, StartSpeech) for cmd in executor.commands)) register_transcribed_turn( runtime._structured_log_context, message_id="long-wait", transcription="fala longa", text="fala longa", ) runtime.on_user_input_transcribed(SimpleNamespace(transcript="fala longa", is_final=True)) await task dispatch.assert_awaited_once() async def test_r2_discards_input_until_its_playout_finishes(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.replay_in_flight = True executor.pipeline_result = BackendReply(stage="PRESENTATION", text="resposta R2") played = [] def _during_r2_playout(_command) -> None: played.append(True) self.assertTrue(runtime.state.deferred_interruption.replay_in_flight) runtime.note_vad_speech_end(1500) runtime.on_user_input_transcribed(SimpleNamespace(transcript="fala ignorada", is_final=True)) executor.wait_for_playout_callbacks.append(_during_r2_playout) with mock.patch.object(runtime, "_play_deferred_short_comfort") as short_comfort: await runtime.run_pipeline("texto R2", "texto R2", user_seq=1, is_deferred_replay=True) self.assertTrue(played, "o playout de R2 precisa acontecer de verdade") short_comfort.assert_not_called() self.assertFalse(runtime.state.deferred_interruption.replay_in_flight) self.assertEqual(runtime.state.user_final_seq, 1) self.assertEqual(runtime.state.deferred_interruption.long_turns, []) async def test_r2_reply_is_not_interruptible(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.current_stage = "ARGUMENTATION" runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.replay_in_flight = True executor.pipeline_result = BackendReply(stage="ARGUMENTATION", text="resposta R2") await runtime.run_pipeline("texto R2", "texto R2", user_seq=1, is_deferred_replay=True) speech_commands = [cmd for cmd in executor.commands if isinstance(cmd, StartSpeech)] self.assertEqual([cmd.text for cmd in speech_commands], ["resposta R2"]) # Descartar a fala no runtime nao basta: com allow_interruptions=True o # barge-in do proprio LiveKit corta o playout e R2 sai pela metade. self.assertFalse(speech_commands[0].allow_interruptions) async def test_r2_logs_and_ignores_short_voice_without_playing_comfort(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.state.deferred_interruption.replay_in_flight = True with ( mock.patch.object(runtime, "_play_deferred_short_comfort") as short_comfort, mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log, ): runtime.note_vad_speech_end(999) short_comfort.assert_not_called() self.assertTrue( any( call.args[1] == "deferred_replay_vad_discarded" and call.kwargs["speech_duration_ms"] == 999 and not call.kwargs["short_comfort_scheduled"] for call in flow_log.call_args_list ) ) async def test_normal_reply_stays_interruptible_outside_the_replay(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.current_stage = "ARGUMENTATION" runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply(stage="ARGUMENTATION", text="resposta normal") await runtime.run_pipeline("texto", "texto", user_seq=1) speech_commands = [cmd for cmd in executor.commands if isinstance(cmd, StartSpeech)] self.assertEqual([cmd.text for cmd in speech_commands], ["resposta normal"]) self.assertTrue(speech_commands[0].allow_interruptions) async def test_r1_waits_for_the_user_turn_to_end_before_deciding(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 backend_started = asyncio.Event() release_backend = asyncio.Event() async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() executor.run_pipeline_side_effect = _wait_for_backend with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort"), mock.patch.object( runtime, "_dispatch_deferred_interruption", new_callable=mock.AsyncMock, return_value=True, ) as dispatch, ): task = asyncio.create_task(runtime.run_pipeline("original", "original", user_seq=1)) await backend_started.wait() # R1 volta com o cliente no meio da fala: nem o VAD fechou, nem ha # STT pendente ainda. runtime._user_not_speaking.clear() release_backend.set() await asyncio.sleep(0) self.assertFalse(task.done()) self.assertFalse(any(isinstance(cmd, StartSpeech) for cmd in executor.commands)) runtime.note_vad_speech_end(1300) runtime._user_not_speaking.set() register_transcribed_turn( runtime._structured_log_context, message_id="late-long", transcription="fala que comecou antes de R1 voltar", text="fala que comecou antes de R1 voltar", ) runtime.on_user_input_transcribed( SimpleNamespace( transcript="fala que comecou antes de R1 voltar", is_final=True ) ) await task dispatch.assert_awaited_once() self.assertFalse(any(isinstance(cmd, StartSpeech) for cmd in executor.commands)) async def test_r1_gives_up_waiting_when_the_stt_final_never_arrives(self) -> None: runtime, _agent, executor = self._make_runtime( runtime_config_overrides={ "deferred_interruption_stt_settle_timeout_s": 0.01, } ) runtime.state.user_final_seq = 1 backend_started = asyncio.Event() release_backend = asyncio.Event() async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() executor.run_pipeline_side_effect = _wait_for_backend executor.pipeline_result = BackendReply(stage="PRESENTATION", text="resposta R1") with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort"), mock.patch.object( runtime, "_speak_backend_reply", new_callable=mock.AsyncMock ) as speak, ): task = asyncio.create_task(runtime.run_pipeline("original", "original", user_seq=1)) await backend_started.wait() runtime.note_vad_speech_end(1200) release_backend.set() # Sem teto na espera esta chamada nunca retornaria. await asyncio.wait_for(task, timeout=5) deferred = runtime.state.deferred_interruption self.assertEqual(deferred.pending_long_stt_finals, 0) self.assertEqual(len(deferred.pending_vad_utterances), 0) self.assertFalse(deferred.backend_in_flight) # A resposta de R1 e liberada em vez de ficar presa a um final que nao vem. speak.assert_awaited_once() async def test_special_comfort_flag_is_cleared_when_r1_answers_empty(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.state.user_final_seq = 1 backend_started = asyncio.Event() release_backend = asyncio.Event() async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() _executor.run_pipeline_side_effect = _wait_for_backend _executor.pipeline_result = BackendReply(stage="PRESENTATION", text="") with mock.patch.object(runtime, "_play_deferred_interruption_comfort"): task = asyncio.create_task(runtime.run_pipeline("original", "original", user_seq=1)) await backend_started.wait() runtime.note_vad_speech_end(1200) runtime.on_user_input_transcribed(SimpleNamespace(transcript="", is_final=True)) release_backend.set() await task # Sem a limpeza aqui o conforto especial nunca mais tocaria na ligacao. self.assertFalse(runtime.state.deferred_interruption.special_comfort_sent) async def test_empty_final_is_counted_once_when_provider_and_livekit_both_report(self) -> None: runtime, _agent, _executor = self._make_runtime() # Em raw_json o transcript cru nao e vazio, so o texto extraido dele e. runtime._extract_text_from_transcript = lambda _transcript: "" deferred = runtime.state.deferred_interruption deferred.backend_in_flight = True runtime.note_vad_speech_end(300) runtime.note_vad_speech_end(1500) with mock.patch.object(runtime, "_play_deferred_interruption_comfort"): # Provider avisa o final vazio da fala curta... runtime.handle_empty_stt_final(source="stt_provider") # ...e o LiveKit repassa o mesmo resultado como evento. runtime.on_user_input_transcribed( SimpleNamespace(transcript='{"data":{"text":""}}', is_final=True) ) await self._drain_tasks() # A fala longa continua na fila, aguardando o proprio final. self.assertEqual(len(deferred.pending_vad_utterances), 1) self.assertEqual(deferred.pending_long_stt_finals, 1) async def test_vad_utterance_without_final_does_not_leak_into_the_next_turn(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 backend_started = asyncio.Event() release_backend = asyncio.Event() async def _wait_for_backend() -> None: backend_started.set() await release_backend.wait() executor.run_pipeline_side_effect = _wait_for_backend with ( mock.patch.object(runtime, "_play_deferred_interruption_comfort"), mock.patch.object(runtime, "_speak_backend_reply", new_callable=mock.AsyncMock), ): task = asyncio.create_task(runtime.run_pipeline("original", "original", user_seq=1)) await backend_started.wait() # Ruido curto que o STT nunca finaliza. runtime.note_vad_speech_end(200) release_backend.set() await asyncio.wait_for(task, timeout=5) deferred = runtime.state.deferred_interruption self.assertEqual(len(deferred.pending_vad_utterances), 0) # No turno seguinte a fala do cliente nao pode ser consumida pela sobra. deferred.backend_in_flight = True with mock.patch.object(runtime, "_play_deferred_interruption_comfort"): runtime.note_vad_speech_end(1500) register_transcribed_turn( runtime._structured_log_context, message_id="turno-seguinte", transcription="agora eu quero falar", text="agora eu quero falar", ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="agora eu quero falar", is_final=True) ) await self._drain_tasks() self.assertEqual( [turn.text for turn in deferred.long_turns], ["agora eu quero falar"] ) async def test_protected_speech_drop_wins_over_the_deferred_cycle(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.state.user_final_seq = 1 deferred = runtime.state.deferred_interruption deferred.backend_in_flight = True runtime._mark_drop_next_user_final( {"reason": "protected_speech", "stage": "PRESENTATION"} ) with mock.patch.object(runtime, "_play_deferred_interruption_comfort"): runtime.note_vad_speech_end(1500) register_transcribed_turn( runtime._structured_log_context, message_id="protegida", transcription="fala descartada", text="fala descartada", ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="fala descartada", is_final=True) ) await self._drain_tasks() self.assertEqual(deferred.long_turns, []) self.assertEqual(runtime.state.user_final_seq, 1) # A fala VAD tem de sair da fila junto, senao o proximo final consumiria # a entrada errada. self.assertEqual(len(deferred.pending_vad_utterances), 0) self.assertEqual(deferred.pending_long_stt_finals, 0) async def test_special_comfort_defers_the_periodic_backend_wait_notice(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.state.deferred_interruption.backend_in_flight = True runtime._inflight_backend_activity_at = 0.0 with mock.patch.object(runtime, "_play_deferred_interruption_comfort"): runtime.note_vad_speech_end(1500) self.assertGreater(runtime._inflight_backend_activity_at, 0.0) def test_combined_interruption_transcription_concatenates_text_and_drops_word_timestamps(self) -> None: turns = [ SimpleNamespace( transcription='{"data":{"text":"primeira", "words":[{"word":"primeira"}]}}', text="primeira", ), SimpleNamespace(transcription="segunda", text="segunda"), ] transcription, text = CallRuntime._combined_interruption_transcription(turns) payload = __import__("json").loads(transcription) self.assertEqual(text, "primeira. segunda") self.assertEqual(payload["data"]["text"], "primeira. segunda") self.assertNotIn("words", payload["data"]) async def test_on_user_input_transcribed_drops_during_protected_feedback(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime.state.current_stage = "PRESENTATION" runtime.state.current_speech.stage = "PRESENTATION" runtime.state.current_speech.drop_user_input_while_speaking = True runtime.state.current_speech.agent_message_type = "feedback" runtime.speaking.set() with mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log: runtime.on_user_input_transcribed( SimpleNamespace(transcript="quero falar agora", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertIsNone(runtime.state.pending_user_final) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( call.args[1] == "user_input_dropped" and call.kwargs["agent_message_type"] == "feedback" and call.kwargs["text"] == "quero falar agora" for call in flow_log.call_args_list ) ) self.assertTrue( any(event == "user_input_dropped" for event, _fields in runtime._timeline.events) ) async def test_dropped_user_final_does_not_leak_message_id_to_repeated_text(self) -> None: runtime, _agent, executor = self._make_runtime() structured_context = StructuredLogContext( callid="call-1", session_id="session-drop-repeated-text", num_telefone="5511999999999", cod_ani="1234", nome_agente="conta", ) runtime._structured_log_context = structured_context reset_turn_message_sequence(structured_context, clear_pending=True) register_started_turn_message_id(structured_context, "MSG-drop-0001") register_transcribed_turn( structured_context, message_id="MSG-drop-0001", transcription="mesmo texto", text="mesmo texto", ) runtime._drop_next_user_final_context = { "reason": "test_drop", "stage": "PRESENTATION", "agent_message_type": "ready", "agent_result_type": "", } try: runtime.on_user_input_transcribed( SimpleNamespace(transcript="mesmo texto", is_final=True) ) await self._drain_tasks() self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) register_started_turn_message_id(structured_context, "MSG-current-0002") register_transcribed_turn( structured_context, message_id="MSG-current-0002", transcription="mesmo texto", text="mesmo texto", ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="mesmo texto", is_final=True) ) await self._drain_tasks() pipeline_command = next(cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput)) self.assertEqual(pipeline_command.user_input["message_id"], "MSG-current-0002") finally: reset_turn_message_sequence(structured_context, clear_pending=True) async def test_on_user_input_transcribed_drops_after_feedback_while_backend_is_processing(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} await runtime._speak_backend_reply( BackendReply( stage="PRESENTATION", text="Ainda estou consultando sua fatura.", metadata={ "event": "feedback", "agent_message_type": "feedback", "expects_user_response": False, "drop_user_input_while_speaking": True, "is_interruptible": False, }, ), source="backend_push", add_to_chat_ctx=True, ) self.assertFalse(runtime.speaking.is_set()) runtime.on_user_input_transcribed( SimpleNamespace(transcript="posso falar?", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( event == "user_input_dropped" and fields["agent_message_type"] == "feedback" for event, fields in runtime._timeline.events ) ) async def test_on_user_input_transcribed_drops_after_result_feedback_while_backend_is_processing(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} await runtime._speak_backend_reply( BackendReply( stage="PRESENTATION", text="Aguarde um instante, por favor.", metadata={ "agent_message_type": "result", "agent_result_type": "feedback", "expects_user_response": False, "drop_user_input_while_speaking": True, "is_interruptible": False, }, ), source="run_pipeline", add_to_chat_ctx=True, ) self.assertFalse(runtime.speaking.is_set()) runtime.on_user_input_transcribed( SimpleNamespace(transcript="voce esta cancelando?", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertIsNone(runtime.state.pending_user_final) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( event == "user_input_dropped" and fields["agent_message_type"] == "result" and fields["agent_result_type"] == "feedback" for event, fields in runtime._timeline.events ) ) async def test_on_user_input_transcribed_drops_during_terminal_speech(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "DONE" runtime.state.current_speech.stage = "DONE" runtime.state.current_speech.drop_user_input_while_speaking = True runtime.state.current_speech.agent_message_type = "result" runtime.state.current_speech.agent_result_type = "resolvido" runtime.speaking.set() runtime.on_user_input_transcribed( SimpleNamespace(transcript="tenho outra pergunta", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertIsNone(runtime.state.pending_user_final) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_user_state_speaking_during_ready_drops_late_final_but_arms_timeout(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Posso verificar sua fatura.", metadata={ "agent_message_type": "ready", "expects_user_response": True, "drop_user_input_while_speaking": True, "is_interruptible": False, "wait_timeout_seconds": 60.0, }, ) executor.wait_for_playout_delay_s = 0.03 runtime.register_callbacks() task = asyncio.create_task(runtime.run_pipeline("", "", user_seq=1)) try: while not any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands): await asyncio.sleep(0) runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) runtime.on_user_input_transcribed( SimpleNamespace(transcript="ja estou respondendo", is_final=True) ) await task pipeline_commands = [cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput)] self.assertEqual(len(pipeline_commands), 1) self.assertEqual(runtime.state.user_final_seq, 1) self.assertTrue( any( event == "user_input_dropped" and fields["agent_message_type"] == "ready" for event, fields in runtime._timeline.events ) ) self.assertTrue( any(event == "agent_wait_timeout_armed" for event, _fields in runtime._timeline.events) ) finally: if not task.done(): task.cancel() await asyncio.gather(task, return_exceptions=True) await self._cancel_pending_tasks() async def test_user_state_speaking_during_noninterruptible_presentation_drops_late_final(self) -> None: runtime, _agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} _agent.pipeline = _FakeInflightPushPipeline() runtime.state.current_stage = "PRESENTATION" runtime.state.current_speech.stage = "PRESENTATION" runtime.state.current_speech.allow_interruptions = False runtime.speaking.set() runtime.register_callbacks() runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") runtime.speaking.clear() runtime.state.current_speech.stage = "" runtime.state.current_speech.allow_interruptions = False runtime.on_user_input_transcribed( SimpleNamespace(transcript="oi", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertIsNone(runtime.state.pending_user_final) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertFalse(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) self.assertTrue( any( event == "pre_backend_wait_notice_skipped" and fields["skipped_reason"] == "user_input_drop_pending" for event, fields in runtime._timeline.events ) ) self.assertTrue( any( event == "user_input_dropped" and fields["agent_message_type"] == "ready" and fields["text"] == "oi" for event, fields in runtime._timeline.events ) ) async def test_user_state_speaking_in_presentation_post_playout_grace_drops_late_final(self) -> None: runtime, _agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} _agent.pipeline = _FakeInflightPushPipeline() runtime.register_callbacks() await runtime.say_stage( "Ola, sou a assistente virtual da TIM.", "PRESENTATION", allow_interruptions=False, schedule_idle_after=False, ) self.assertFalse(runtime.speaking.is_set()) runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") runtime.on_user_input_transcribed( SimpleNamespace(transcript="estou falando junto", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 0) self.assertIsNone(runtime.state.pending_user_final) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( event == "user_input_drop_grace_armed" and fields["grace_ms"] == 250 for event, fields in runtime._timeline.events ) ) self.assertTrue( any( event == "pre_backend_wait_notice_skipped" and fields["skipped_reason"] == "user_input_drop_pending" for event, fields in runtime._timeline.events ) ) self.assertTrue( any( event == "user_input_dropped" and fields["reason"] == "user_state_speaking_protected_post_playout_grace" and fields["agent_message_type"] == "ready" and fields["text"] == "estou falando junto" for event, fields in runtime._timeline.events ) ) async def test_presentation_post_playout_grace_expires_before_next_user_final(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} with mock.patch( "app.livekit.runtime.call_runtime.PROTECTED_SPEECH_POST_PLAYOUT_DROP_GRACE_S", 0.001, ): await runtime.say_stage( "Ola, sou a assistente virtual da TIM.", "PRESENTATION", allow_interruptions=False, schedule_idle_after=False, ) await asyncio.sleep(0.01) runtime.on_user_input_transcribed( SimpleNamespace(transcript="agora posso responder", is_final=True) ) await self._drain_tasks() self.assertEqual(runtime.state.user_final_seq, 1) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_user_state_speaking_interrupts_current_tts(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "ARGUMENTATION" runtime.state.current_speech.stage = "ARGUMENTATION" runtime.state.current_speech.allow_interruptions = True runtime.state.current_speech.handle = executor.speech_handle runtime.state.current_speech.speech_id = "7f3a0000-0000-4000-8000-000000000000" runtime.speaking.set() runtime.register_callbacks() runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) await self._drain_tasks() interrupt_commands = [ cmd for cmd in executor.commands if isinstance(cmd, InterruptSpeech) ] self.assertEqual(len(interrupt_commands), 1) self.assertIs(interrupt_commands[0].handle, executor.speech_handle) async def test_user_state_speaking_interrupts_inflight_wait_notice(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "AGENT_BACKEND_WAIT" runtime.state.current_speech.stage = "AGENT_BACKEND_WAIT" runtime.state.current_speech.allow_interruptions = False runtime.state.current_speech.handle = executor.speech_handle runtime.speaking.set() runtime.register_callbacks() runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) await self._drain_tasks() interrupt_commands = [ cmd for cmd in executor.commands if isinstance(cmd, InterruptSpeech) ] self.assertEqual(len(interrupt_commands), 1) self.assertIs(interrupt_commands[0].handle, executor.speech_handle) self.assertTrue(interrupt_commands[0].force) async def test_user_state_speaking_cancels_agent_wait_timeout(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.register_callbacks() runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Ola, como posso ajudar?", metadata={"wait_timeout_seconds": 0.01}, ), user_seq=0, source="test", ) runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="listening", new_state="speaking") ) await asyncio.sleep(0.03) await self._drain_tasks() self.assertFalse(any(isinstance(cmd, NotifyBridgeStop) for cmd in executor.commands)) self.assertFalse(runtime.finalized.is_set()) async def test_user_state_listening_does_not_speak_wait_notice_before_stt_final(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.register_callbacks() runtime._session.handlers["user_state_changed"]( SimpleNamespace(old_state="speaking", new_state="listening") ) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, []) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_vad_wait_notice_starts_without_waiting_for_user_state_listening(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.register_callbacks() runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) async def test_vad_wait_notice_fast_path_uses_vad_pause_without_user_state(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, "pre_backend_wait_notice_fast_on_vad_pause": True, } ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} runtime._user_not_speaking.clear() agent.pipeline = _FakeInflightPushPipeline() with mock.patch( "app.livekit.runtime.call_runtime.asyncio.sleep", new=mock.AsyncMock(), ) as sleep: runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() sleep.assert_not_awaited() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) self.assertTrue( any( event == "pre_backend_wait_notice_started" and fields["fast_on_vad_pause"] is True for event, fields in runtime._timeline.events ) ) async def test_vad_wait_notice_default_path_skips_while_user_state_speaking(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} runtime._user_not_speaking.clear() agent.pipeline = _FakeInflightPushPipeline() with mock.patch("app.livekit.runtime.call_runtime.PRE_BACKEND_WAIT_NOTICE_GUARD_S", 0): runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() self.assertFalse( any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands) ) self.assertTrue( any( event == "pre_backend_wait_notice_skipped" and fields["skipped_reason"] == "user_speaking" for event, fields in runtime._timeline.events ) ) def _wait_notice_config(self, wait_text: str) -> dict: return { "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } async def test_vad_wait_notice_skipped_while_agent_is_speaking(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides=self._wait_notice_config(wait_text) ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.state.current_stage = "PRESENTATION" runtime.state.current_speech.stage = "PRESENTATION" runtime.speaking.set() runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() self.assertFalse( any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands) ) self.assertTrue( any( event == "pre_backend_wait_notice_skipped" and fields["skipped_reason"] == "agent_speaking" for event, fields in runtime._timeline.events ) ) self.assertFalse(runtime._pre_backend_wait_notice_reserved) async def test_vad_wait_notice_skipped_while_backend_is_processing(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides=self._wait_notice_config(wait_text) ) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() await agent._run_lock.acquire() try: runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() finally: agent._run_lock.release() self.assertFalse( any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands) ) self.assertTrue( any( event == "pre_backend_wait_notice_skipped" and fields["skipped_reason"] == "backend_processing" for event, fields in runtime._timeline.events ) ) self.assertFalse(runtime._pre_backend_wait_notice_reserved) async def test_wait_notice_audio_aborts_when_another_speech_took_the_lock(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides=self._wait_notice_config(wait_text) ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() say_stage_seq_at_schedule = runtime._say_stage_seq await runtime.say_stage("Resposta do agente.", "PRESENTATION") played = await runtime._play_inflight_backend_wait_notice_audio( user_seq=None, attempt=1, abort_if_agent_spoke_since=say_stage_seq_at_schedule, ) await self._drain_tasks() self.assertFalse(played) speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, ["Resposta do agente."]) async def test_vad_wait_notice_does_not_queue_behind_an_agent_speech(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides=self._wait_notice_config(wait_text) ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() # O aviso e agendado no fim de fala do VAD e so pega o say_lock depois # que a fala do agente libera: nesse ponto ele nao deve mais sair. runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await runtime.say_stage("Resposta do agente.", "PRESENTATION") await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, ["Resposta do agente."]) async def test_vad_wait_notice_returns_reservation_when_it_does_not_play(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides=self._wait_notice_config(wait_text) ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await runtime.say_stage("Resposta do agente.", "PRESENTATION") await self._drain_tasks() self.assertFalse(runtime._pre_backend_wait_notice_reserved) # A reserva devolvida nao pode bloquear o proximo agendamento. runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual( speech_texts, [ "Resposta do agente.", _expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text), ], ) async def test_vad_wait_notice_uses_started_turn_message_id(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) structured_context = StructuredLogContext( callid="call-1", session_id="session-started-turn-runtime", num_telefone="5511999999999", cod_ani="1234", nome_agente="conta", ) runtime._structured_log_context = structured_context runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() reset_turn_message_sequence(structured_context, clear_pending=True) register_started_turn_message_id(structured_context, "GED-test-0002") try: runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() wait_events = [ fields for event, fields in runtime._timeline.events if event == "inflight_backend_wait_audio" ] self.assertEqual(len(wait_events), 1) self.assertTrue(wait_events[0]["message_id"].startswith("GED-test-0002_conforto_")) finally: reset_turn_message_sequence(structured_context, clear_pending=True) async def test_vad_wait_notice_uses_last_agent_message_id_without_started_turn(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) structured_context = StructuredLogContext( callid="call-1", session_id="session-last-agent-turn-runtime", num_telefone="5511999999999", cod_ani="1234", nome_agente="conta", ) runtime._structured_log_context = structured_context runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() reset_turn_message_sequence(structured_context, clear_pending=True) runtime._remember_agent_base_message_id("MSG-agent-0001") try: runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") await self._drain_tasks() wait_events = [ fields for event, fields in runtime._timeline.events if event == "inflight_backend_wait_audio" ] self.assertEqual(len(wait_events), 1) self.assertTrue(wait_events[0]["message_id"].startswith("MSG-agent-0001_conforto_")) finally: reset_turn_message_sequence(structured_context, clear_pending=True) async def test_inflight_backend_wait_notice_skips_post_user_final_grace(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, _executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.last_user_final_at = time.monotonic() agent.pipeline = _FakeInflightPushPipeline() with mock.patch.object( runtime, "_wait_post_user_final_grace", new=mock.AsyncMock(), ) as wait_grace: await runtime._play_inflight_backend_wait_notice_audio( user_seq=None, attempt=1, ) wait_grace.assert_not_awaited() async def test_empty_stt_final_logs_and_clears_pre_backend_wait_notice(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime._pre_backend_wait_notice_reserved = True with ( mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log, mock.patch.object(runtime, "arm_idle_timer") as arm_idle, ): runtime.on_user_input_transcribed( SimpleNamespace(transcript="", is_final=True) ) await self._drain_tasks() self.assertFalse(runtime._pre_backend_wait_notice_reserved) arm_idle.assert_called_once_with(reason="empty_stt_final") self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertTrue( any( call.args[1] == "stt_final_empty" and call.kwargs["reason"] == "empty_text" and call.kwargs["pre_backend_wait_reserved"] is True for call in flow_log.call_args_list ) ) self.assertIn( ( "user_transcript_final_empty", { "current_stage": "INTRO", "speaking": False, "transcript_len": 0, "pre_backend_wait_reserved": True, }, ), runtime._timeline.events, ) async def test_empty_stt_final_rearms_agent_wait_timeout_when_available(self) -> None: runtime, _agent, _executor = self._make_runtime() reply = BackendReply( stage="PRESENTATION", text="Ainda esta ai?", metadata={"wait_timeout_seconds": 60.0}, ) runtime.arm_agent_wait_timeout_from_reply(reply, user_seq=0, source="test") runtime.cancel_agent_wait_timeout("user_state_speaking") try: with ( mock.patch.object( runtime, "_arm_agent_wait_timeout", wraps=runtime._arm_agent_wait_timeout, ) as arm_wait, mock.patch.object(runtime, "arm_idle_timer") as arm_idle, ): runtime.on_user_input_transcribed( SimpleNamespace(transcript="", is_final=True) ) arm_wait.assert_called_once() self.assertEqual(arm_wait.call_args.kwargs["user_seq"], 0) self.assertEqual(arm_wait.call_args.kwargs["source"], "empty_stt_final") arm_idle.assert_not_called() finally: await self._cancel_pending_tasks() async def test_empty_stt_final_during_comfort_rearms_after_playout(self) -> None: runtime, _agent, executor = self._make_runtime() reply = BackendReply( stage="PRESENTATION", text="Ainda esta ai?", metadata={"wait_timeout_seconds": 60.0}, ) runtime.arm_agent_wait_timeout_from_reply(reply, user_seq=0, source="test") runtime.cancel_agent_wait_timeout("user_state_speaking") executor.wait_for_playout_callbacks.append( lambda _command: runtime.handle_empty_stt_final(source="stt_provider") ) try: with mock.patch.object( runtime, "_arm_agent_wait_timeout", wraps=runtime._arm_agent_wait_timeout, ) as arm_wait: await runtime.say_stage( "Um instantinho", "AGENT_BACKEND_WAIT", add_to_chat_ctx=False, allow_interruptions=False, schedule_idle_after=False, ) arm_wait.assert_called_once() self.assertEqual(arm_wait.call_args.kwargs["user_seq"], 0) self.assertEqual(arm_wait.call_args.kwargs["source"], "empty_stt_final") self.assertFalse(runtime._empty_stt_recovery_pending) finally: await self._cancel_pending_tasks() async def test_inflight_backend_wait_audio_leaves_tts_metrics_empty_without_exporting_original_id( self, ) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._timeline = _FakeTimeline() agent.pipeline = _FakeInflightPushPipeline() executor.speech_handle = SimpleNamespace(id="speech-no-metric", interrupted=False) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch.dict(os.environ, {"TTS_TTFB_METRIC_WAIT_MS": "0"}), mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event, ): await runtime._play_inflight_backend_wait_notice_audio( user_seq=1, attempt=2, message_id="GED-test-0002", ) self.assertEqual(len(structured_events), 1) self.assertIsNone(structured_events[0]["latencia_total_ms"]) self.assertIsNone(structured_events[0]["latencia_tffb_ms"]) self.assertIsNone(structured_events[0]["duracao_audio_ms"]) self.assertNotIn("original_message_id", structured_events[0]) wait_events = [ fields for event, fields in runtime._timeline.events if event == "inflight_backend_wait_audio" ] self.assertEqual(len(wait_events), 1) self.assertTrue(wait_events[0]["message_id"].startswith("GED-test-0002_conforto_")) self.assertEqual(wait_events[0]["text"], _expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)) self.assertNotIn("original_message_id", wait_events[0]) async def test_inflight_backend_wait_audio_uses_sidecar_text_for_logged_event(self) -> None: with TemporaryDirectory() as tmp_dir: long_dir = Path(tmp_dir) / "long" long_dir.mkdir() audio_path = long_dir / "01.wav" text_path = long_dir / "01.txt" audio_path.write_bytes(b"fake-wav") text_path.write_text( "Estou verificando as informacoes para te ajudar.\nSo um momentinho.", encoding="utf-8", ) expected_text = "Estou verificando as informacoes para te ajudar. So um momentinho." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": "Texto padrao que nao deve aparecer.", "inflight_backend_wait_long_audio_dir": str(long_dir), } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._timeline = _FakeTimeline() agent.pipeline = _FakeInflightPushPipeline() with ( mock.patch("app.livekit.runtime.call_runtime.wav_duration_ms", return_value=100), mock.patch("app.livekit.runtime.call_runtime.wav_audio_frames", return_value="audio-frames"), ): await runtime._play_inflight_backend_wait_notice_audio( user_seq=1, attempt=2, message_id="GED-test-0002", ) speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual([cmd.text for cmd in speech_commands], [expected_text]) wait_events = [ fields for event, fields in runtime._timeline.events if event == "inflight_backend_wait_audio" ] self.assertEqual(len(wait_events), 1) self.assertEqual(wait_events[0]["text"], expected_text) self.assertEqual(wait_events[0]["path"], str(audio_path)) async def test_backend_wait_notice_starts_immediately_after_stt_final(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 0.01, "inflight_backend_wait_max_notices": 1, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.state.user_final_seq = 1 never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() executor.run_pipeline_side_effect = _hang_backend await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertTrue(speech_commands[0].allow_interruptions) self.assertIsNotNone(speech_commands[0].audio) async def test_pre_stt_wait_notice_activity_controls_next_backend_notice_interval(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 0.01, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() runtime._pre_backend_wait_notice_reserved = True runtime._inflight_backend_activity_at = time.monotonic() - 120.0 never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() executor.run_pipeline_side_effect = _hang_backend await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual( [cmd.text for cmd in speech_commands], [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)], ) self.assertTrue(speech_commands[0].allow_interruptions) async def test_inflight_backend_wait_notice_skips_when_user_is_speaking(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime._user_not_speaking.clear() await runtime._play_inflight_backend_wait_notice_audio( user_seq=1, attempt=1, message_id="GED-test-0001", ) await self._drain_tasks() self.assertFalse( any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands) ) async def test_pre_stt_wait_notice_is_not_restarted_while_pending(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 300.0, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} agent.pipeline = _FakeInflightPushPipeline() runtime.schedule_pre_backend_wait_notice(reason="vad_end_of_speech") runtime.schedule_pre_backend_wait_notice(reason="user_state_listening") await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) async def test_stt_final_does_not_interrupt_pre_backend_wait_notice(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._pre_backend_wait_notice_reserved = True runtime.speaking.set() runtime.state.current_stage = "PRESENTATION" runtime.state.current_speech.stage = "AGENT_BACKEND_WAIT" runtime.state.current_speech.allow_interruptions = True runtime.state.current_speech.handle = executor.speech_handle runtime.on_user_input_transcribed( SimpleNamespace(transcript="Por que veio mais caro?", is_final=True) ) await self._drain_tasks() self.assertFalse(any(isinstance(cmd, InterruptSpeech) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_sim_during_pre_backend_wait_notice_runs_pipeline(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._pre_backend_wait_notice_reserved = True runtime.speaking.set() runtime.state.current_stage = "PRESENTATION" runtime.state.current_speech.stage = "AGENT_BACKEND_WAIT" runtime.state.current_speech.allow_interruptions = False runtime.state.current_speech.handle = executor.speech_handle with mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log: runtime.on_user_input_transcribed( SimpleNamespace(transcript="sim", is_final=True) ) await self._drain_tasks() self.assertFalse(any(isinstance(cmd, InterruptSpeech) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertFalse( any( call.args[1] == "interrupt_ignored" and call.kwargs.get("reason") == "backchannel" and call.kwargs.get("text") == "sim" for call in flow_log.call_args_list ) ) async def test_run_pipeline_forwards_pending_speech_interruption(self) -> None: runtime, agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 agent._consumed_interrupt = ( "Claro, vou te explicar", "7f3a0000-0000-4000-8000-000000000000", False, ) await runtime.run_pipeline("texto", "texto", user_seq=1) interruption_commands = [ cmd for cmd in executor.commands if isinstance(cmd, SetPipelineInterruption) ] self.assertEqual(len(interruption_commands), 1) self.assertTrue(interruption_commands[0].interrupted) self.assertEqual(interruption_commands[0].listened_text, "Claro, vou te explicar") self.assertEqual( interruption_commands[0].speech_id, "7f3a0000-0000-4000-8000-000000000000", ) async def test_deferred_replay_marks_processing_interruption(self) -> None: runtime, agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 agent._consumed_interrupt = ( "", "7f3a0000-0000-4000-8000-000000000001", True, ) await runtime.run_pipeline( "texto interrompido", "texto interrompido", user_seq=1, is_deferred_replay=True, ) command = next( cmd for cmd in executor.commands if isinstance(cmd, SetPipelineProcessingInterruption) ) self.assertEqual(command.listened_text, "") self.assertEqual( command.speech_id, "7f3a0000-0000-4000-8000-000000000001", ) async def test_run_pipeline_uses_registered_turn_message_id(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 structured_context = StructuredLogContext( callid="call-1", session_id="session-1", num_telefone="5511999999999", cod_ani="1234", nome_agente="conta", ) runtime._structured_log_context = structured_context reset_turn_message_sequence(structured_context, clear_pending=True) register_transcribed_turn( structured_context, message_id="MSG-call-1-0001", transcription="texto", text="texto", ) await runtime.run_pipeline("texto", "texto", user_seq=1) pipeline_command = next(cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput)) self.assertEqual(pipeline_command.user_input["message_id"], "MSG-call-1-0001") enriched_reply = runtime._backend_reply_with_message_id( BackendReply(stage="PRESENTATION", text="resposta"), "MSG-call-1-0001", ) self.assertEqual(runtime._message_id_from_reply(enriched_reply), "MSG-call-1-0001") async def test_on_user_input_transcribed_clears_started_turn_message_id(self) -> None: runtime, _agent, executor = self._make_runtime() structured_context = StructuredLogContext( callid="call-1", session_id="session-user-final-started-turn", num_telefone="5511999999999", cod_ani="1234", nome_agente="conta", ) runtime._structured_log_context = structured_context reset_turn_message_sequence(structured_context, clear_pending=True) register_started_turn_message_id(structured_context, "MSG-call-1-0002") register_transcribed_turn( structured_context, message_id="MSG-call-1-0002", transcription="texto", text="texto", ) try: runtime.on_user_input_transcribed( SimpleNamespace(transcript="texto", is_final=True) ) await self._drain_tasks() pipeline_command = next(cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput)) self.assertEqual(pipeline_command.user_input["message_id"], "MSG-call-1-0002") self.assertEqual(peek_started_turn_message_id(structured_context), "") finally: reset_turn_message_sequence(structured_context, clear_pending=True) async def test_say_stage_marks_interruption_with_speech_id(self) -> None: runtime, _agent, executor = self._make_runtime() executor.speech_handle = SimpleNamespace(interrupted=True) executor.spoken_text = "Claro, vou te explicar" await runtime.say_stage( "Claro, vou te explicar sua fatura.", "ARGUMENTATION", allow_interruptions=True, speech_id="7f3a0000-0000-4000-8000-000000000000", ) pending_commands = [cmd for cmd in executor.commands if isinstance(cmd, SetPendingInterrupt)] self.assertEqual(len(pending_commands), 1) self.assertEqual(pending_commands[0].listened_text, "Claro, vou te explicar") self.assertEqual( pending_commands[0].speech_id, "7f3a0000-0000-4000-8000-000000000000", ) async def test_backend_reply_metadata_controls_tts_interruptibility(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply( stage="ARGUMENTATION", text="Ola, como posso ajudar?", metadata={ "speech_id": "11111111-1111-4111-8111-111111111111", "is_interruptible": False, }, ) await runtime.run_pipeline("texto", "texto", user_seq=1) speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] self.assertEqual(len(speech_commands), 1) self.assertFalse(speech_commands[0].allow_interruptions) async def test_run_pipeline_skips_stale_turn_before_running_pipeline(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 2 await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) async def test_run_pipeline_forwards_short_acknowledgement_while_speaking(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.speaking.set() with mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log: await runtime.run_pipeline("ok", "ok", user_seq=0) self.assertTrue(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertFalse( any( call.args[1] == "interrupt_ignored" and call.kwargs.get("reason") == "backchannel" for call in flow_log.call_args_list ) ) async def test_low_confidence_sim_after_tts_runs_pipeline(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} await runtime.say_stage( "Consegui esclarecer sua dúvida?", "PRESENTATION", allow_interruptions=True, ) self.assertFalse(runtime.speaking.is_set()) low_confidence_sim = stt_text_with_single_word_threshold( { "data": { "text": "sim", "words": [{"word": "sim", "probability": 0.009}], } }, min_prob_single_word=0.03, ) self.assertEqual(low_confidence_sim, "sim") with mock.patch("app.livekit.runtime.call_runtime.log_flow_event") as flow_log: runtime.on_user_input_transcribed( SimpleNamespace(transcript=low_confidence_sim, is_final=True) ) await self._drain_tasks() pipeline_commands = [ cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput) ] self.assertTrue(pipeline_commands) self.assertEqual(pipeline_commands[-1].user_input, "sim") self.assertFalse( any( call.args[1] == "interrupt_ignored" and call.kwargs.get("reason") == "backchannel" for call in flow_log.call_args_list ) ) async def test_run_pipeline_drops_output_if_user_speaks_during_pipeline(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 async def _mutate_seq(): runtime.state.user_final_seq = 2 executor.pipeline_result = BackendReply(stage="PRESENTATION", text="resposta") executor.run_pipeline_side_effect = _mutate_seq await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertEqual(runtime.state.current_stage, "INTRO") self.assertFalse(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) async def test_say_stage_waits_briefly_after_user_final(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.last_user_final_at = time.monotonic() sleep_calls = [] async def _fake_sleep(delay: float) -> None: sleep_calls.append(delay) with mock.patch("app.livekit.runtime.call_runtime.asyncio.sleep", side_effect=_fake_sleep): await runtime.say_stage("resposta", "ARGUMENTATION") self.assertTrue(sleep_calls) self.assertGreater(sleep_calls[0], 0.0) self.assertLessEqual(sleep_calls[0], 0.35) self.assertTrue(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) async def test_say_stage_logs_tts_metrics_and_success_status(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._record_xai_tts_turn_timing( { "segment_id": "segment-test", "max_audio_delta_gap_ms": 840, "max_playout_underrun_0ms": 125, "xai_micro_underflows": 2, "xai_avg_underrun_ms": 80, } ) runtime._record_tts_metric( SimpleNamespace( metrics=SimpleNamespace( type="tts_metrics", ttfb=0.123, duration=0.456, audio_duration=1.500, speech_id="speech-test", segment_id="segment-test", label="provider", ) ) ) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event): await runtime.say_stage("resposta", "ARGUMENTATION", message_id="GED-555-0002") self.assertTrue(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) self.assertEqual(len(structured_events), 1) event = structured_events[0] self.assertEqual(event["tipo_evento"], "envio msg") self.assertIsNotNone(event["inicio_ns"]) self.assertIsNotNone(event["fim_ns"]) self.assertEqual(event["latencia_total_ms"], 456) self.assertEqual(event["latencia_tffb_ms"], 123) self.assertEqual(event["duracao_audio_ms"], 1500) self.assertEqual(event["tts_max_gap_ms"], 840) self.assertEqual(event["tts_max_underrun_0ms"], 125) self.assertEqual(event["tts_underflow_count"], 2) self.assertEqual(event["tts_avg_underflow_ms"], 80) self.assertEqual(event["message_id"], "GED-555-0002") self.assertEqual(event["http_cod_status"], 200) self.assertEqual(event["http_cod_desc"], "OK") async def test_say_stage_leaves_tts_metrics_empty_when_metric_is_missing(self) -> None: runtime, _agent, executor = self._make_runtime() executor.speech_handle = SimpleNamespace(id="speech-without-metric", interrupted=False) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch.dict(os.environ, {"TTS_TTFB_METRIC_WAIT_MS": "0"}), mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event, ): await runtime.say_stage("resposta", "ARGUMENTATION", message_id="GED-555-0003") self.assertEqual(len(structured_events), 1) self.assertIsNone(structured_events[0]["latencia_total_ms"]) self.assertIsNone(structured_events[0]["latencia_tffb_ms"]) self.assertIsNone(structured_events[0]["duracao_audio_ms"]) async def test_say_stage_logs_audio_duration_from_tts_metrics(self) -> None: runtime, _agent, executor = self._make_runtime() executor.speech_handle = SimpleNamespace(id="speech-with-audio-metric", interrupted=False) runtime._record_tts_metric( SimpleNamespace( metrics=SimpleNamespace( type="tts_metrics", ttfb=0.123, duration=0.400, audio_duration=1.500, speech_id="speech-with-audio-metric", label="provider", ) ) ) structured_events = [] flow_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) def _capture_flow_event(_logger, event_name, **kwargs): flow_events.append((event_name, kwargs)) with mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event, ), mock.patch( "app.livekit.runtime.call_runtime.log_flow_event", side_effect=_capture_flow_event, ): await runtime.say_stage("resposta", "ARGUMENTATION", message_id="GED-555-0004") self.assertEqual(len(structured_events), 1) self.assertEqual(structured_events[0]["latencia_total_ms"], 400) self.assertEqual(structured_events[0]["latencia_tffb_ms"], 123) self.assertEqual(structured_events[0]["duracao_audio_ms"], 1500) tts_done_events = [fields for event, fields in flow_events if event == "tts_done"] self.assertEqual(tts_done_events[0]["duration_ms"], 400) self.assertEqual(tts_done_events[0]["provider_duration_ms"], 400) self.assertEqual(tts_done_events[0]["metrics_source"], "livekit_tts_metrics") async def test_say_stage_marks_interruption_with_cancelled_tts_metric(self) -> None: runtime, _agent, executor = self._make_runtime() executor.speech_handle = SimpleNamespace(id="speech-cancelled", interrupted=True) executor.spoken_text = "preciso falar" runtime._record_tts_metric( SimpleNamespace( metrics=SimpleNamespace( type="tts_metrics", ttfb=0.150, duration=0.700, audio_duration=0.250, speech_id="speech-cancelled", label="provider", cancelled=True, ) ) ) structured_events = [] flow_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) def _capture_flow_event(_logger, event_name, **kwargs): flow_events.append((event_name, kwargs)) with mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event, ), mock.patch( "app.livekit.runtime.call_runtime.log_flow_event", side_effect=_capture_flow_event, ): await runtime.say_stage( "resposta interrompida", "ARGUMENTATION", message_id="GED-555-0005", allow_interruptions=True, ) self.assertEqual(len(structured_events), 1) self.assertEqual(structured_events[0]["latencia_total_ms"], 700) self.assertEqual(structured_events[0]["latencia_tffb_ms"], 150) self.assertEqual(structured_events[0]["duracao_audio_ms"], 250) self.assertIs(structured_events[0]["interrupcao"], True) tts_done_events = [fields for event, fields in flow_events if event == "tts_done"] self.assertIs(tts_done_events[0]["interruption"], True) pending_commands = [cmd for cmd in executor.commands if isinstance(cmd, SetPendingInterrupt)] self.assertEqual(len(pending_commands), 1) self.assertEqual(pending_commands[0].listened_text, "preciso falar") async def test_say_stage_keeps_interruption_duration_without_first_tts_audio(self) -> None: runtime, _agent, executor = self._make_runtime() executor.speech_handle = SimpleNamespace(id="speech-cancelled-no-audio", interrupted=True) runtime._record_tts_metric( SimpleNamespace( metrics=SimpleNamespace( type="tts_metrics", ttfb=-1.0, duration=0.700, audio_duration=0.0, speech_id="speech-cancelled-no-audio", label="provider", cancelled=True, ) ) ) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event, ): await runtime.say_stage( "resposta interrompida", "ARGUMENTATION", message_id="GED-555-0006", allow_interruptions=True, ) self.assertEqual(len(structured_events), 1) self.assertEqual(structured_events[0]["latencia_total_ms"], 700) self.assertIsNone(structured_events[0]["latencia_tffb_ms"]) self.assertIsNone(structured_events[0]["duracao_audio_ms"]) self.assertIs(structured_events[0]["interrupcao"], True) async def test_say_stage_returns_to_listening_when_no_first_audio_is_seen(self) -> None: runtime, _agent, executor = self._make_runtime() first_handle = SimpleNamespace(id="speech-first", interrupted=False) executor.speech_handles = [first_handle] executor.wait_for_playout_delays_s = [0.05] flow_events = [] def capture(_logger, event_name, **kwargs): flow_events.append((event_name, kwargs)) with mock.patch.dict(os.environ, {"TTS_PLAYOUT_START_TIMEOUT_S": "0.01", "TTS_TTFB_METRIC_WAIT_MS": "0"}), mock.patch("app.livekit.runtime.call_runtime.log_flow_event", side_effect=capture): result = await runtime.say_stage("resposta", "ARGUMENTATION", message_id="GED-555-0005") self.assertFalse(result) speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] self.assertEqual(len(speech_commands), 1) self.assertFalse(runtime.finalized.is_set()) self.assertTrue(first_handle.interrupted) self.assertTrue(any(event == "tts_failure_return_to_listening" for event, _ in flow_events)) async def test_say_stage_returns_to_listening_after_partial_audio_failure(self) -> None: runtime, _agent, executor = self._make_runtime() first_handle = SimpleNamespace(id="speech-gap", interrupted=False) executor.speech_handles = [first_handle] executor.wait_for_playout_exceptions = [RuntimeError("xAI TTS partial audio failure: underflow_error; socket_resynchronized=0; xai_underrun_estimado_ms=1001; xai_micro_underflows=1; xai_avg_underrun_ms=1001; underflow_error_ms=1000; max_audio_delta_gap_ms=1200; pcm_duration_ms=500; pcm_bytes=24000")] structured_events = [] def capture_event(*_args, **kwargs): structured_events.append(kwargs) with mock.patch.dict(os.environ, {"TTS_TTFB_METRIC_WAIT_MS": "0"}), mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=capture_event): result = await runtime.say_stage("resposta com falha no meio", "ARGUMENTATION", message_id="GED-555-0006") self.assertFalse(result) speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] self.assertEqual(len(speech_commands), 1) self.assertFalse(runtime.finalized.is_set()) self.assertTrue(first_handle.interrupted) self.assertEqual(structured_events[0]["erro_msg"], "Falha TTS") self.assertIn("underflow_error", structured_events[0]["erro_detalhe"]) self.assertIn("underflow_error_ms=1000", structured_events[0]["erro_detalhe"]) self.assertEqual(structured_events[0]["tts_max_gap_ms"], 1200) self.assertEqual(structured_events[0]["tts_max_underrun_0ms"], 1001) self.assertEqual(structured_events[0]["tts_underflow_count"], 1) self.assertEqual(structured_events[0]["tts_avg_underflow_ms"], 1001) async def test_say_stage_publishes_partial_failure_without_livekit_metrics(self) -> None: runtime, _agent, executor = self._make_runtime() first_handle = SimpleNamespace(id="speech-partial-event", interrupted=False) executor.speech_handles = [first_handle] structured_events = [] def emit_partial_failure(_command) -> None: runtime._record_xai_tts_turn_failed( { "segment_id": "segment-partial-event", "reason": "underflow_error", "max_audio_delta_gap_ms": 920, "max_playout_underrun_0ms": 1001, "xai_micro_underflows": 3, "xai_avg_underrun_ms": 611, "provider_synthesis_ms": 2450, "pcm_duration_ms": 860, } ) executor.wait_for_playout_callbacks = [emit_partial_failure] def capture_event(*_args, **kwargs): structured_events.append(kwargs) with mock.patch.dict(os.environ, {"TTS_TTFB_METRIC_WAIT_MS": "0"}), mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=capture_event, ): result = await runtime.say_stage( "resposta com falha parcial", "ARGUMENTATION", message_id="GED-555-0008", ) self.assertFalse(result) self.assertTrue(first_handle.interrupted) self.assertEqual(len(structured_events), 1) event = structured_events[0] self.assertEqual(event["erro_msg"], "Falha TTS") self.assertIn("underflow_error", event["erro_detalhe"]) self.assertEqual(event["latencia_total_ms"], 2450) self.assertEqual(event["duracao_audio_ms"], 860) self.assertEqual(event["tts_max_gap_ms"], 920) self.assertEqual(event["tts_max_underrun_0ms"], 1001) self.assertEqual(event["tts_underflow_count"], 3) self.assertEqual(event["tts_avg_underflow_ms"], 611) self.assertIsNone(event["http_cod_status"]) async def test_say_stage_returns_to_listening_after_swallowed_tts_error(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.register_callbacks() first_handle = SimpleNamespace(id="speech-swallowed", interrupted=False) executor.speech_handles = [first_handle] def emit_error(_command): error = SimpleNamespace(type="tts_error", error=RuntimeError("no audio frames were pushed for text: resposta"), recoverable=False) runtime._session.handlers["error"](SimpleNamespace(error=error)) executor.wait_for_playout_callbacks = [emit_error] with mock.patch.dict(os.environ, {"TTS_TTFB_METRIC_WAIT_MS": "0"}): result = await runtime.say_stage("resposta", "PRESENTATION", message_id="GED-555-0007") self.assertFalse(result) speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] self.assertEqual(len(speech_commands), 1) self.assertFalse(runtime.finalized.is_set()) async def test_tts_metric_zero_is_not_used_as_ttfb(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._record_tts_metric( SimpleNamespace( metrics=SimpleNamespace( type="tts_metrics", ttfb=0, speech_id="speech-zero", label="provider", ) ) ) self.assertNotIn("speech-zero", runtime._tts_tffb_ms_by_speech_handle_id) async def test_finalize_is_idempotent(self) -> None: runtime, _agent, executor = self._make_runtime() await asyncio.gather(runtime.finalize("first"), runtime.finalize("second")) self.assertTrue(runtime.finalized.is_set()) self.assertEqual(sum(isinstance(cmd, EndServiceOnce) for cmd in executor.commands), 1) self.assertEqual(sum(isinstance(cmd, ExportSession) for cmd in executor.commands), 1) async def test_run_pipeline_done_triggers_done_notification_and_finalize(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply(stage="DONE", text="encerrando", done=True) await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() self.assertTrue(any(isinstance(cmd, NotifyBridgeDone) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_run_pipeline_non_terminal_reply_does_not_finalize_from_global_done_stage( self, ) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Perfeito! Seguiremos com o cancelamento.", done=False, ) def _mark_global_stage_done(_command) -> None: runtime.state.current_stage = "DONE" executor.wait_for_playout_callbacks = [_mark_global_stage_done] await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() self.assertFalse(any(isinstance(cmd, NotifyBridgeDone) for cmd in executor.commands)) self.assertFalse(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertFalse(runtime.finalized.is_set()) async def test_run_pipeline_agent_final_result_speaks_before_terminal_stop(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply( stage="DONE", text="Atendimento encerrado com sucesso.", done=True, export_payload={ "type": "nao_resolvido", "content": "Atendimento encerrado com sucesso.", "tool_calls": [], }, ) executor.end_reply = BackendReply( stage="DONE", text="Nao deveria falar de novo.", done=True, export_payload=[{"status": "ok"}], ) await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertEqual([cmd.text for cmd in speech_commands], ["Atendimento encerrado com sucesso."]) self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_nao_resolvido") self.assertEqual(stop_commands[0].reason, "nao_resolvido") self.assertEqual(stop_commands[0].phase, "in_session") self.assertFalse(any(isinstance(cmd, NotifyBridgeDone) for cmd in executor.commands)) self.assertLess( executor.commands.index(speech_commands[0]), executor.commands.index(stop_commands[0]), ) self.assertTrue(runtime.finalized.is_set()) self.assertFalse(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertFalse(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_terminal_reply_waits_until_user_stops_speaking(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._user_not_speaking.clear() executor.pipeline_result = BackendReply( stage="DONE", text="Atendimento encerrado.", done=True, export_payload={"type": "nao_resolvido", "content": "Atendimento encerrado."}, ) with mock.patch( "app.livekit.runtime.call_runtime.POST_USER_FINAL_SPEECH_GRACE_S", 0.01, ): task = asyncio.create_task(runtime.run_pipeline("texto", "texto", user_seq=1)) await asyncio.sleep(0.01) self.assertFalse(any(isinstance(cmd, StartSpeech) for cmd in executor.commands)) runtime._user_not_speaking.set() await task self.assertEqual( [cmd.text for cmd in executor.commands if isinstance(cmd, StartSpeech)], ["Atendimento encerrado."], ) async def test_terminal_reply_is_spoken_when_speaking_becomes_new_user_turn(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._user_not_speaking.clear() executor.pipeline_result = BackendReply( stage="DONE", text="Atendimento encerrado.", done=True, export_payload={"type": "nao_resolvido", "content": "Atendimento encerrado."}, ) with mock.patch( "app.livekit.runtime.call_runtime.POST_USER_FINAL_SPEECH_GRACE_S", 0.01, ): task = asyncio.create_task(runtime.run_pipeline("texto", "texto", user_seq=1)) await asyncio.sleep(0.01) runtime._user_not_speaking.set() runtime.state.user_final_seq = 2 await task speech_commands = [cmd for cmd in executor.commands if isinstance(cmd, StartSpeech)] self.assertEqual([cmd.text for cmd in speech_commands], ["Atendimento encerrado."]) self.assertFalse(speech_commands[0].allow_interruptions) self.assertTrue(any(isinstance(cmd, NotifyBridgeStop) for cmd in executor.commands)) self.assertTrue(runtime.finalized.is_set()) async def test_run_pipeline_skips_new_turn_after_done_stage(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.current_stage = "DONE" runtime.state.user_final_seq = 2 await runtime.run_pipeline("Ok,", "Ok,", user_seq=2) self.assertFalse(any(isinstance(cmd, RunPipelineInput) for cmd in executor.commands)) self.assertFalse(any(isinstance(cmd, StartSpeech) for cmd in executor.commands)) async def test_run_pipeline_ready_wait_timeout_sends_long_silence_stop(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 0 executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Ola, como posso ajudar com sua fatura?", metadata={"wait_timeout_seconds": 0.01}, ) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event): await runtime.run_pipeline("", "", user_seq=0) await self._drain_tasks() speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertEqual([cmd.text for cmd in speech_commands], ["Ola, como posso ajudar com sua fatura?"]) self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_silencio_longo") self.assertEqual(stop_commands[0].reason, "no_user_response") self.assertEqual(stop_commands[0].phase, "in_session") self.assertTrue( any( event["tipo_evento"] == "recebimento msg" and event.get("erro_msg") == "Silencio Longo" for event in structured_events ) ) self.assertTrue(runtime.finalized.is_set()) self.assertFalse(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertFalse(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_run_pipeline_ready_wait_timeout_sends_retry_messages_before_stop(self) -> None: controller = _FakeVADThresholdController() runtime, _agent, executor = self._make_runtime(vad_threshold_controller=controller) runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 0 retry_messages = [ "Voce ainda esta na linha?", "Sigo por aqui.", "Vou confirmar mais uma vez.", ] executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Ola, como posso ajudar com sua fatura?", metadata={ "message_id": "MSG-timeout-0001", "wait_timeout_seconds": 0.01, "wait_retry_messages": retry_messages, }, ) structured_events = [] def _capture_structured_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch( "app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_structured_event, ): await runtime.run_pipeline("", "", user_seq=0) await self._drain_tasks() speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertEqual( [cmd.text for cmd in speech_commands], ["Ola, como posso ajudar com sua fatura?", *retry_messages], ) self.assertEqual( [cmd.add_to_chat_ctx for cmd in speech_commands], [True, False, False, False], ) self.assertEqual( [cmd.allow_interruptions for cmd in speech_commands[1:]], [True, True, True], ) self.assertIn(("activate", 1, "agent_wait_timeout_retry"), controller.calls) self.assertIn(("restore", "agent_wait_timeout:no_user_response"), controller.calls) event_ids = [event["message_id"] for event in structured_events] self.assertEqual( event_ids[:4], [ "MSG-timeout-0001", "MSG-timeout-0001_inatividade_1", "MSG-timeout-0001_inatividade_2", "MSG-timeout-0001_inatividade_3", ], ) self.assertEqual(len(event_ids), 5) self.assertIn("_erro_terminal_recebimento_", event_ids[4]) retry_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_retry" ] self.assertEqual( [event["message_id"] for event in retry_events], [ "MSG-timeout-0001_inatividade_1", "MSG-timeout-0001_inatividade_2", "MSG-timeout-0001_inatividade_3", ], ) self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_silencio_longo") self.assertEqual(stop_commands[0].reason, "no_user_response") self.assertEqual(stop_commands[0].phase, "in_session") self.assertLess( executor.commands.index(speech_commands[-1]), executor.commands.index(stop_commands[0]), ) self.assertTrue(runtime.finalized.is_set()) self.assertFalse(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertFalse(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_agent_wait_timeout_retry_registers_idle_nudge_event(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._timeline = _FakeTimeline() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 0 retry_messages = [ "Voce ainda esta na linha?", "Sigo por aqui.", "Vou confirmar mais uma vez.", ] executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Ola, como posso ajudar com sua fatura?", metadata={ "message_id": "MSG-timeout-0001", "wait_timeout_seconds": 0.01, "wait_retry_messages": retry_messages, }, ) await runtime.run_pipeline("", "", user_seq=0) await self._drain_tasks() nudge_commands = [cmd for cmd in executor.commands if isinstance(cmd, InjectIdleNudge)] self.assertEqual([cmd.text for cmd in nudge_commands], retry_messages) # O evento e bufferizado ANTES da fala, para sobreviver a uma interrupcao # do cliente em cima do retry. speech_commands = [cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech"] for nudge, speech in zip(nudge_commands, speech_commands[1:]): self.assertLess( executor.commands.index(nudge), executor.commands.index(speech), ) async def test_agent_wait_timeout_metadata_persists_for_later_replies(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._timeline = _FakeTimeline() try: runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Ola, como posso ajudar?", metadata={ "wait_timeout_seconds": 60.0, "wait_retry_messages": ["Voce ainda esta na linha?"], }, ), user_seq=0, source="initial", ) runtime.cancel_agent_wait_timeout("user_final") runtime._clear_agent_wait_timeout_memory() runtime.state.user_final_seq = 1 runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Pode me responder?", ), user_seq=1, source="next_reply", ) armed_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_armed" ] self.assertEqual(len(armed_events), 2) self.assertEqual(armed_events[-1]["wait_timeout_s"], 60.0) self.assertEqual(armed_events[-1]["retry_messages"], 1) self.assertEqual(armed_events[-1]["source"], "next_reply") self.assertEqual(armed_events[-1]["user_seq"], 1) finally: await self._cancel_pending_tasks() async def test_agent_wait_timeout_metadata_rewrites_persisted_config(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._timeline = _FakeTimeline() try: runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Primeira pergunta", metadata={ "wait_timeout_seconds": 60.0, "wait_retry_messages": ["Mensagem antiga"], }, ), user_seq=0, source="initial", ) runtime.cancel_agent_wait_timeout("user_final") runtime._clear_agent_wait_timeout_memory() runtime.state.user_final_seq = 1 runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Nova pergunta", metadata={ "wait_timeout_seconds": 30.0, "wait_retry_messages": ["Nova mensagem", "Ultimo aviso"], }, ), user_seq=1, source="rewrite", ) runtime.cancel_agent_wait_timeout("user_final") runtime._clear_agent_wait_timeout_memory() runtime.state.user_final_seq = 2 runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Pergunta sem metadata", ), user_seq=2, source="after_rewrite", ) armed_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_armed" ] self.assertEqual(len(armed_events), 3) self.assertEqual(armed_events[-2]["wait_timeout_s"], 30.0) self.assertEqual(armed_events[-2]["retry_messages"], 2) self.assertEqual(armed_events[-1]["wait_timeout_s"], 30.0) self.assertEqual(armed_events[-1]["retry_messages"], 2) self.assertEqual(armed_events[-1]["source"], "after_rewrite") self.assertEqual(armed_events[-1]["user_seq"], 2) finally: await self._cancel_pending_tasks() async def test_feedback_reply_does_not_arm_client_wait_timeout(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._timeline = _FakeTimeline() try: runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Posso verificar sua fatura.", metadata={ "agent_message_type": "ready", "expects_user_response": True, "drop_user_input_while_speaking": True, "is_interruptible": False, "wait_timeout_seconds": 60.0, "wait_retry_messages": ["Voce esta ai?"], }, ), user_seq=0, source="ready", ) runtime.cancel_agent_wait_timeout("test") runtime._clear_agent_wait_timeout_memory() runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Ainda estou consultando sua fatura.", metadata={ "agent_message_type": "feedback", "expects_user_response": False, "drop_user_input_while_speaking": True, "is_interruptible": False, "event": "feedback", }, ), user_seq=0, source="feedback", ) armed_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_armed" ] self.assertEqual(len(armed_events), 1) runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Consegui consultar. Posso te ajudar em mais algo?", metadata={ "agent_message_type": "result", "agent_result_type": "final", "expects_user_response": True, "drop_user_input_while_speaking": False, }, ), user_seq=0, source="final", ) armed_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_armed" ] self.assertEqual(len(armed_events), 2) self.assertEqual(armed_events[-1]["source"], "final") self.assertEqual(armed_events[-1]["retry_messages"], 1) finally: await self._cancel_pending_tasks() async def test_agent_wait_timeout_zero_metadata_clears_persisted_config(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime._timeline = _FakeTimeline() try: runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Primeira pergunta", metadata={"wait_timeout_seconds": 60.0}, ), user_seq=0, source="initial", ) runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Sem timeout daqui em diante", metadata={"wait_timeout_seconds": 0}, ), user_seq=0, source="disable", ) runtime.state.user_final_seq = 1 runtime.arm_agent_wait_timeout_from_reply( BackendReply( stage="PRESENTATION", text="Pergunta sem metadata", ), user_seq=1, source="after_disable", ) armed_events = [ fields for event, fields in runtime._timeline.events if event == "agent_wait_timeout_armed" ] self.assertEqual(len(armed_events), 1) finally: await self._cancel_pending_tasks() async def test_run_pipeline_backend_failure_triggers_terminal_stop_and_finalize(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 async def _fail_backend() -> None: raise RuntimeError("backend closed connection") executor.run_pipeline_side_effect = _fail_backend structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event): await runtime.run_pipeline("texto", "texto", user_seq=1) stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_agent_backend_unavailable") self.assertEqual(stop_commands[0].reason, "resource_unhealthy") self.assertEqual(stop_commands[0].resource, "agent_backend") self.assertEqual(stop_commands[0].failed_resources, ("agent_backend",)) self.assertTrue( any( event["tipo_evento"] == "envio msg" and event.get("erro_msg") == "Falha comunicacao" and "_erro_terminal_envio_" in event.get("message_id", "") for event in structured_events ) ) self.assertTrue(any(isinstance(cmd, EndServiceOnce) for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_run_pipeline_transfer_reply_uses_suffixed_message_id(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.state.user_final_seq = 1 executor.pipeline_result = BackendReply( stage="PRESENTATION", text="", done=False, export_payload={ "type": "transferred", "status": "transferred", "additionalInformations": {"changeAgentId": "ATH-001"}, }, ) structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event): await runtime.run_pipeline( "texto", "texto", user_seq=1, message_id="GED-test-0002", ) self.assertTrue( any( event["tipo_evento"] == "envio msg" and event.get("erro_msg") == "Transferido" and event.get("message_id", "").startswith("GED-test-0002_transferencia_") for event in structured_events ) ) async def test_finalize_speaks_end_reply_when_remote_listener_exists(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} executor.end_reply = BackendReply( stage="DONE", text="Mensagem final do backend.", done=True, export_payload=[{"status": "ok"}], ) await runtime.finalize("manual") self.assertTrue(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_finalize_skips_end_reply_tts_without_remote_listener(self) -> None: runtime, _agent, executor = self._make_runtime() executor.end_reply = BackendReply( stage="DONE", text="Mensagem final do backend.", done=True, export_payload=[{"status": "ok"}], ) await runtime.finalize("manual") self.assertFalse(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) self.assertTrue(any(isinstance(cmd, ExportSession) for cmd in executor.commands)) async def test_register_callbacks_wires_session_and_shutdown(self) -> None: runtime, _agent, _executor = self._make_runtime() runtime.register_callbacks() self.assertIn("user_input_transcribed", runtime._session.handlers) self.assertIn("close", runtime._session.handlers) self.assertIn("data_received", runtime._ctx.room.handlers) self.assertEqual(len(runtime._ctx.shutdown_callbacks), 1) async def test_session_close_error_triggers_resource_stop(self) -> None: runtime, _agent, executor = self._make_runtime() runtime.register_callbacks() runtime._session.handlers["close"]( SimpleNamespace( reason=SimpleNamespace(value="error"), error=SimpleNamespace(type="stt_error"), ) ) await self._drain_tasks() stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_stt_unavailable") self.assertEqual(stop_commands[0].resource, "stt") async def test_bridge_control_client_audio_enabled_triggers_initial_agent_turn(self) -> None: runtime, _agent, executor = self._make_runtime(agent_starts_conversation=True) runtime._ctx.room.remote_participants = {"user-1": object()} executor.pipeline_result = BackendReply(stage="PRESENTATION", text="Ola, sou o agent.") runtime.register_callbacks() runtime._ctx.room.handlers["data_received"]( SimpleNamespace( topic="bridge.control", data=b'{"type":"client_audio_enabled","room":"room-test","protocol":"PRT-1"}', ) ) await self._drain_tasks() run_commands = [cmd for cmd in executor.commands if isinstance(cmd, RunPipelineInput)] self.assertEqual(len(run_commands), 1) self.assertEqual(run_commands[0].user_input, "") self.assertTrue(any(cmd.__class__.__name__ == "StartSpeech" for cmd in executor.commands)) async def test_run_holds_user_audio_input_before_starting_session(self) -> None: runtime, _agent, executor = self._make_runtime(agent_starts_conversation=True) runtime._ctx.room.remote_participants = {"user-1": object()} audio_enabled_on_start = [] executor.on_start_session = lambda _cmd: audio_enabled_on_start.append( runtime._session.input.audio_enabled ) await runtime.run() await asyncio.sleep(0) await self._cancel_pending_tasks() # O RoomIO le o estado do input ao anexar a track: se o hold viesse depois # do StartSession, o audio do setup ja teria chegado no VAD/STT. self.assertEqual(audio_enabled_on_start, [False]) self.assertFalse(runtime._session.input.audio_enabled) self.assertEqual(runtime._session.input.calls, [False]) async def test_run_keeps_user_audio_input_when_user_starts_conversation(self) -> None: runtime, _agent, _executor = self._make_runtime(agent_starts_conversation=False) runtime._ctx.room.remote_participants = {"user-1": object()} await runtime.run() await asyncio.sleep(0) await self._cancel_pending_tasks() self.assertTrue(runtime._session.input.audio_enabled) self.assertEqual(runtime._session.input.calls, []) async def test_initial_agent_turn_releases_user_audio_input_after_playout(self) -> None: runtime, _agent, executor = self._make_runtime(agent_starts_conversation=True) runtime._ctx.room.remote_participants = {"user-1": object()} executor.pipeline_result = BackendReply(stage="PRESENTATION", text="Ola, sou o agent.") audio_enabled_during_playout = [] executor.wait_for_playout_callbacks = [ lambda _cmd: audio_enabled_during_playout.append( runtime._session.input.audio_enabled ) ] runtime.hold_user_audio_input(reason="room_setup") runtime.register_callbacks() runtime._ctx.room.handlers["data_received"]( SimpleNamespace( topic="bridge.control", data=b'{"type":"client_audio_enabled","room":"room-test","protocol":"PRT-1"}', ) ) await self._drain_tasks() self.assertEqual(audio_enabled_during_playout, [False]) self.assertTrue(runtime._session.input.audio_enabled) self.assertEqual(runtime._session.input.calls, [False, True]) async def test_initial_agent_turn_releases_user_audio_input_when_pipeline_fails(self) -> None: runtime, _agent, executor = self._make_runtime(agent_starts_conversation=True) runtime._ctx.room.remote_participants = {"user-1": object()} async def _boom() -> None: raise RuntimeError("backend fora do ar") executor.run_pipeline_side_effect = _boom runtime.hold_user_audio_input(reason="room_setup") runtime.register_callbacks() runtime._ctx.room.handlers["data_received"]( SimpleNamespace( topic="bridge.control", data=b'{"type":"client_audio_enabled","room":"room-test","protocol":"PRT-1"}', ) ) await self._drain_tasks() # Falha no turno inicial nao pode deixar a chamada surda. self.assertTrue(runtime._session.input.audio_enabled) async def test_user_audio_input_gate_timeout_releases_when_initial_turn_never_starts( self, ) -> None: runtime, _agent, _executor = self._make_runtime(agent_starts_conversation=True) runtime._ctx.room.remote_participants = {"user-1": object()} with mock.patch.dict(os.environ, {"USER_AUDIO_INPUT_SETUP_TIMEOUT_S": "0.01"}): runtime.hold_user_audio_input(reason="room_setup") runtime._arm_user_audio_input_gate_timeout() await self._drain_tasks() # client_audio_enabled nunca chegou: o gate abre sozinho. self.assertTrue(runtime._session.input.audio_enabled) self.assertEqual(runtime._session.input.calls, [False, True]) async def test_backend_push_message_is_spoken_without_new_user_turn(self) -> None: runtime, _agent, executor = self._make_runtime( agent_starts_conversation=True, push_replies=[BackendReply(stage="PRESENTATION", text="Segunda mensagem do backend.")], ) runtime._ctx.room.remote_participants = {"user-1": object()} executor.pipeline_result = BackendReply(stage="PRESENTATION", text="Primeira mensagem do backend.") runtime.register_callbacks() runtime._ctx.room.handlers["data_received"]( SimpleNamespace( topic="bridge.control", data=b'{"type":"client_audio_enabled","room":"room-test","protocol":"PRT-1"}', ) ) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual( speech_texts, [ "Primeira mensagem do backend.", "Segunda mensagem do backend.", ], ) async def test_inflight_backend_feedback_is_spoken_before_final_reply(self) -> None: runtime, agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 pipeline = _FakeInflightPushPipeline() agent.pipeline = pipeline executor.pipeline_result = BackendReply( stage="PRESENTATION", text="Resposta final do backend.", ) async def _emit_feedback_then_finish() -> None: await pipeline.push( BackendReply( stage="PRESENTATION", text="Ainda estou consultando sua fatura.", metadata={ "event": "feedback", "agent_message_type": "feedback", "expects_user_response": False, "drop_user_input_while_speaking": True, "is_interruptible": False, }, ) ) await asyncio.sleep(0.01) pipeline.close() executor.run_pipeline_side_effect = _emit_feedback_then_finish await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual( speech_texts, [ "Ainda estou consultando sua fatura.", "Resposta final do backend.", ], ) speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertFalse(speech_commands[0].allow_interruptions) async def test_inflight_backend_wait_notice_is_spoken_immediately_after_stt_final(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 0.01, "inflight_backend_wait_max_notices": 1, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() executor.run_pipeline_side_effect = _hang_backend await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) async def test_inflight_backend_wait_audio_sequence_starts_short_then_long(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." with TemporaryDirectory() as tmp_dir: base_dir = Path(tmp_dir) short_dir = base_dir / "short" long_dir = base_dir / "long" short_dir.mkdir() long_dir.mkdir() short_audio = short_dir / "curto.wav" long_audio = long_dir / "longo.wav" short_audio.write_bytes(b"short") long_audio.write_bytes(b"long") runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.005, "inflight_backend_wait_timeout_s": 0.6, "inflight_backend_wait_max_notices": 2, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_short_audio_dir": str(short_dir), "inflight_backend_wait_long_audio_dir": str(long_dir), } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._timeline = _FakeTimeline() agent.pipeline = _FakeInflightPushPipeline() never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() def _audio_frames(path: str): return f"audio:{Path(path).name}" executor.run_pipeline_side_effect = _hang_backend with ( mock.patch("app.livekit.runtime.call_runtime.wav_duration_ms", return_value=100), mock.patch( "app.livekit.runtime.call_runtime.wav_audio_frames", side_effect=_audio_frames, ), mock.patch.object(runtime, "_log_resource_error_events") as structured_error, ): await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual( [cmd.audio for cmd in speech_commands], ["audio:curto.wav", "audio:longo.wav"], ) max_notice_events = [ fields for event, fields in runtime._timeline.events if event == "inflight_backend_wait_max_notices_reached" ] self.assertEqual(len(max_notice_events), 1) self.assertEqual(max_notice_events[0]["notices_sent"], 2) self.assertEqual(max_notice_events[0]["max_notices"], 2) max_notice_structured_calls = [ call for call in structured_error.call_args_list if call.kwargs.get("reason") == "inflight_backend_wait_max_notices_reached" ] self.assertEqual(len(max_notice_structured_calls), 1) self.assertEqual( max_notice_structured_calls[0].kwargs["resource"], "agent_backend", ) async def test_inflight_backend_wait_zero_max_notices_is_unlimited(self) -> None: runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.005, "inflight_backend_wait_timeout_s": 0.2, "inflight_backend_wait_max_notices": 0, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._timeline = _FakeTimeline() agent.pipeline = _FakeInflightPushPipeline() backend_done = asyncio.Event() attempts: list[int] = [] async def _hang_backend() -> None: await backend_done.wait() async def _play_notice(**kwargs): attempts.append(kwargs["attempt"]) if len(attempts) == 3: backend_done.set() return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ): await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertEqual(attempts, [1, 2, 3]) self.assertFalse( any( event == "inflight_backend_wait_max_notices_reached" for event, _fields in runtime._timeline.events ) ) async def test_inflight_backend_wait_skipped_notice_does_not_consume_attempt(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." with TemporaryDirectory() as tmp_dir: base_dir = Path(tmp_dir) short_dir = base_dir / "short" short_dir.mkdir() short_audio = short_dir / "curto.wav" short_audio.write_bytes(b"short") runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.005, "inflight_backend_wait_timeout_s": 0.06, "inflight_backend_wait_max_notices": 1, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_short_audio_dir": str(short_dir), } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime._user_not_speaking.clear() agent.pipeline = _FakeInflightPushPipeline() never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() async def _mark_user_not_speaking() -> None: await asyncio.sleep(0.015) runtime._user_not_speaking.set() def _audio_frames(path: str): return f"audio:{Path(path).name}" executor.run_pipeline_side_effect = _hang_backend restore_task = asyncio.create_task(_mark_user_not_speaking()) try: with ( mock.patch("app.livekit.runtime.call_runtime.wav_duration_ms", return_value=100), mock.patch( "app.livekit.runtime.call_runtime.wav_audio_frames", side_effect=_audio_frames, ), ): await runtime.run_pipeline("texto", "texto", user_seq=1) finally: restore_task.cancel() await asyncio.gather(restore_task, return_exceptions=True) await self._drain_tasks() speech_commands = [ cmd for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual([cmd.audio for cmd in speech_commands], ["audio:curto.wav"]) async def test_inflight_backend_wait_interval_counts_after_audio_playout(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.025, "inflight_backend_wait_timeout_s": 0.04, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() executor.wait_for_playout_delay_s = 0.02 never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() executor.run_pipeline_side_effect = _hang_backend await runtime.run_pipeline("texto", "texto", user_seq=1) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] self.assertEqual(speech_texts, [_expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text)]) async def test_short_processing_interruption_rearms_periodic_notice_deadline(self) -> None: interval_s = 0.05 runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": interval_s, "inflight_backend_wait_timeout_s": 0.3, "inflight_backend_wait_max_notices": 2, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() backend_done = asyncio.Event() attempts: list[int] = [] attempted_at: list[float] = [] async def _hang_backend() -> None: await backend_done.wait() async def _finish_short_speech() -> None: await asyncio.sleep(0.01) runtime.note_vad_speech_end(300) runtime._user_not_speaking.set() async def _play_notice(**kwargs): attempt = kwargs["attempt"] attempts.append(attempt) attempted_at.append(time.monotonic()) if attempts == [1]: return InflightBackendWaitNoticeResult(started=True) if attempts == [1, 2]: runtime._user_not_speaking.clear() asyncio.create_task(_finish_short_speech()) return InflightBackendWaitNoticeResult( started=True, interrupted_by_user=True, ) backend_done.set() return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ): await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertEqual(attempts, [1, 2, 2]) self.assertGreaterEqual(attempted_at[2] - attempted_at[1], interval_s) async def test_short_speech_during_silent_backend_wait_does_not_rearm_notice_deadline(self) -> None: interval_s = 0.05 runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": interval_s, "inflight_backend_wait_timeout_s": 0.3, "inflight_backend_wait_max_notices": 2, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() backend_done = asyncio.Event() attempts: list[int] = [] attempted_at: list[float] = [] short_speech_ended_at = 0.0 async def _hang_backend() -> None: await backend_done.wait() async def _emit_short_speech_during_wait() -> None: nonlocal short_speech_ended_at await asyncio.sleep(interval_s / 2) runtime.note_vad_speech_end(300) short_speech_ended_at = time.monotonic() async def _play_notice(**kwargs): attempts.append(kwargs["attempt"]) attempted_at.append(time.monotonic()) if attempts == [1]: asyncio.create_task(_emit_short_speech_during_wait()) else: backend_done.set() return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ): await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertEqual(attempts, [1, 2]) self.assertGreater(short_speech_ended_at, 0.0) self.assertLess(attempted_at[1] - short_speech_ended_at, interval_s) async def test_valid_processing_interruption_rearms_periodic_notice_deadline(self) -> None: interval_s = 0.04 runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": interval_s, "inflight_backend_wait_timeout_s": 0.3, "inflight_backend_wait_max_notices": 2, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() backend_done = asyncio.Event() attempts: list[int] = [] attempted_at: list[float] = [] async def _hang_backend() -> None: await backend_done.wait() async def _finish_valid_speech() -> None: await asyncio.sleep(0.01) runtime.note_vad_speech_end(1000) runtime._user_not_speaking.set() async def _play_notice(**kwargs): attempt = kwargs["attempt"] attempts.append(attempt) attempted_at.append(time.monotonic()) if attempts == [1]: return InflightBackendWaitNoticeResult(started=True) if attempts == [1, 2]: runtime._user_not_speaking.clear() asyncio.create_task(_finish_valid_speech()) return InflightBackendWaitNoticeResult( started=True, interrupted_by_user=True, ) backend_done.set() return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with ( mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ), mock.patch.object(runtime, "_play_deferred_interruption_comfort"), ): await runtime.run_pipeline("texto", "texto", user_seq=1) self.assertEqual(attempts, [1, 2, 2]) self.assertGreaterEqual( attempted_at[2] - attempted_at[1], interval_s, ) async def test_long_processing_interruption_keeps_periodic_notices_active(self) -> None: runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.01, "inflight_backend_wait_timeout_s": 0.2, "inflight_backend_wait_max_notices": 2, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True agent.pipeline = _FakeInflightPushPipeline() backend_done = asyncio.Event() attempts: list[int] = [] async def _hang_backend() -> None: await backend_done.wait() async def _interrupt_after_first_notice() -> None: await asyncio.sleep(0) runtime.note_vad_speech_end(1500) runtime.on_user_input_transcribed( SimpleNamespace(transcript="e nao resolve", is_final=True) ) async def _play_notice(**kwargs): attempts.append(kwargs["attempt"]) if attempts == [1]: asyncio.create_task(_interrupt_after_first_notice()) else: backend_done.set() return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with ( mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ), mock.patch.object(runtime, "_play_deferred_interruption_comfort"), ): await runtime._execute_pipeline_with_inflight_backend_wait( "pedido original", user_seq=1, ) self.assertEqual(runtime.state.user_final_seq, 2) self.assertEqual( [turn.text for turn in runtime.state.deferred_interruption.long_turns], ["e nao resolve"], ) self.assertEqual(attempts, [1, 2]) async def test_long_processing_interruption_does_not_disable_backend_timeout(self) -> None: runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 60.0, "inflight_backend_wait_timeout_s": 0.03, "inflight_backend_wait_max_notices": 1, "inflight_backend_wait_text": DEFAULT_WAIT_TEXT, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 runtime.state.deferred_interruption.backend_in_flight = True agent.pipeline = _FakeInflightPushPipeline() never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() async def _play_notice(**_kwargs): runtime.note_vad_speech_end(1500) runtime.on_user_input_transcribed( SimpleNamespace(transcript="continua demorando", is_final=True) ) return InflightBackendWaitNoticeResult(started=True) executor.run_pipeline_side_effect = _hang_backend with ( mock.patch.object( runtime, "_play_inflight_backend_wait_notice_audio", side_effect=_play_notice, ), mock.patch.object(runtime, "_play_deferred_interruption_comfort"), ): with self.assertRaises(InflightBackendWaitTimedOut): await runtime._execute_pipeline_with_inflight_backend_wait( "pedido original", user_seq=1, ) self.assertEqual(runtime.state.user_final_seq, 2) self.assertTrue(runtime.finalized.is_set()) async def test_inflight_backend_wait_notice_then_backend_unavailable_stop(self) -> None: wait_text = "Um momento, ainda estou consultando para te ajudar." runtime, agent, executor = self._make_runtime( runtime_config_overrides={ "inflight_backend_wait_interval_s": 0.01, "inflight_backend_wait_timeout_s": 0.035, "inflight_backend_wait_max_notices": 6, "inflight_backend_wait_text": wait_text, "inflight_backend_wait_long_audio_path": WAIT_LONG_AUDIO_PATH, } ) runtime._ctx.room.remote_participants = {"user-1": object()} runtime.state.user_final_seq = 1 agent.pipeline = _FakeInflightPushPipeline() never_finishes = asyncio.Event() async def _hang_backend() -> None: await never_finishes.wait() executor.run_pipeline_side_effect = _hang_backend structured_events = [] def _capture_event(*args, **kwargs): structured_events.append(kwargs) with mock.patch("app.livekit.runtime.call_runtime.log_structured_event", side_effect=_capture_event): await runtime.run_pipeline( "texto", "texto", user_seq=1, message_id="GED-test-0002", ) await self._drain_tasks() speech_texts = [ cmd.text for cmd in executor.commands if cmd.__class__.__name__ == "StartSpeech" ] stop_commands = [cmd for cmd in executor.commands if isinstance(cmd, NotifyBridgeStop)] self.assertGreaterEqual(len(speech_texts), 1) self.assertTrue( all( text == _expected_wait_text_for_audio(WAIT_LONG_AUDIO_PATH, wait_text) for text in speech_texts ) ) comfort_events = [ event for event in structured_events if str(event.get("message_id") or "").startswith("GED-test-0002_conforto_") ] self.assertGreaterEqual(len(comfort_events), 1) self.assertTrue( any( event["tipo_evento"] == "envio msg" and event.get("erro_msg") == "Falha comunicacao" and "_erro_terminal_envio_" in event.get("message_id", "") for event in structured_events ) ) self.assertEqual(len(stop_commands), 1) self.assertEqual(stop_commands[0].status, "stop_agent_backend_unavailable") self.assertEqual(stop_commands[0].reason, "resource_unhealthy") self.assertEqual(stop_commands[0].resource, "agent_backend") self.assertEqual(stop_commands[0].failed_resources, ("agent_backend",)) self.assertTrue(runtime.finalized.is_set()) async def test_run_starts_session_with_bridge_identity(self) -> None: runtime, _agent, executor = self._make_runtime() runtime._ctx.room.remote_participants = {"user-1": object()} await runtime.run() await asyncio.sleep(0) await self._cancel_pending_tasks() start_commands = [cmd for cmd in executor.commands if isinstance(cmd, StartSession)] self.assertEqual(len(start_commands), 1) room_options = start_commands[0].room_options participant_identity = getattr(room_options, "participant_identity", None) if participant_identity is None: participant_identity = room_options.kwargs["participant_identity"] self.assertEqual(participant_identity, "bridge-1") self.assertIn("user_input_transcribed", runtime._session.handlers) if __name__ == "__main__": unittest.main() #só pra fazer um novo PR