from __future__ import annotations import time import tempfile import unittest from pathlib import Path from owner_voice_pet.agent_memory import FaissIndexManifest, FakeMemoryManager, MemoryRecordInput, SQLiteMemoryManager from owner_voice_pet.config import AppConfig from owner_voice_pet.full_duplex_audio import FakeWebRtcAudioProcessingProvider, RenderReferenceRingBuffer from owner_voice_pet.full_duplex_control import CancellationGraph, FullDuplexStateMachine from owner_voice_pet.full_duplex_response import ( FakeStreamingLlmProvider, FakeStreamingTtsProvider, InterruptiblePlaybackQueue, LlmStreamEvent, SentenceSegmenter, ) from owner_voice_pet.full_duplex_speech import FakeStreamingSttProvider, FakeVadProvider, InterruptionDetector, TranscriptEvent from owner_voice_pet.full_duplex_runtime import FullDuplexAgentRuntime from owner_voice_pet.full_duplex_testing import ( PerformanceMetricRecorder, build_fake_full_duplex_audio_fixture, diagnostics_contain_sensitive_data, sanitize_diagnostics, ) from owner_voice_pet.llm import MockLlmProvider from owner_voice_pet.models import AudioFrame, AudioSegment, Message, PipelineState, Transcript from owner_voice_pet.stt import MetadataSttProvider from owner_voice_pet.tool_router import MemorySearchTool, ToolCallRequest, ToolContext, ToolRouter from owner_voice_pet.transport import MemoryAudioTransport from owner_voice_pet.tts import SineTtsProvider from owner_voice_pet.vad import EnergyVadProvider, PrimarySpeakerVadRecorder, VadRecorder class RecordingMetadataSttProvider(MetadataSttProvider): def __init__(self) -> None: super().__init__() self.calls: list[AudioSegment] = [] def transcribe(self, segment: AudioSegment) -> Transcript: self.calls.append(segment) return super().transcribe(segment) class PlaybackInjectedTransport(MemoryAudioTransport): def __init__( self, frames: list[AudioFrame], *, injected_frames: list[AudioFrame], inject_after_play_count: int = 1, ) -> None: super().__init__(frames, flush_clears_input=False) self.injected_frames = list(injected_frames) self.inject_after_play_count = inject_after_play_count self.play_count = 0 def play_pcm(self, segment: AudioSegment, interrupt: bool = False): result = super().play_pcm(segment, interrupt=interrupt) self.play_count += 1 if self.play_count == self.inject_after_play_count: for frame in self.injected_frames: self.inject(frame) time.sleep(0.01) return result class FullDuplexIntegrationTests(unittest.TestCase): def test_run_agent_live_runtime_once_consumes_audio_and_replies(self) -> None: frames = [ AudioFrame(b"\xff\x7f", 16000, 1, 0, 1, {"duration_ms": 20, "speech": True, "transcript": "直接提问"}), AudioFrame(b"\xff\x7f", 16000, 1, 20, 2, {"duration_ms": 20, "speech": True, "transcript": "直接提问"}), AudioFrame(b"\x00\x00", 16000, 1, 40, 3, {"duration_ms": 20, "speech": False, "transcript": "直接提问"}), AudioFrame(b"\x00\x00", 16000, 1, 60, 4, {"duration_ms": 20, "speech": False, "transcript": "直接提问"}), ] transport = MemoryAudioTransport(frames, flush_clears_input=False) runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, memory_enabled=False, tool_router_enabled=False, noise_filter_enabled=False, vad_min_duration_ms=40, vad_end_silence_ms=40, end_chime_enabled=False, ), transport=transport, vad_recorder=VadRecorder(EnergyVadProvider(), min_duration_ms=40, end_silence_ms=40), stt=MetadataSttProvider(), llm=MockLlmProvider(["这是全双工回答。"]), tts=SineTtsProvider(), ) summary = runtime.run(once=True) self.assertEqual(summary.completed_turns, 1) self.assertEqual(summary.failed_turns, 0) self.assertEqual([message.role for message in runtime.context.messages()], ["user", "assistant"]) self.assertEqual(runtime.context.messages()[0].content, "直接提问") self.assertEqual(runtime.context.messages()[1].content, "这是全双工回答。") self.assertGreaterEqual(len(transport.played_segments), 1) self.assertIsNotNone(runtime.audio_hub) self.assertGreater(runtime.audio_hub.processed_capture.frame_count, 0) def test_run_agent_live_uses_primary_speaker_endpoint_instead_of_max_recording(self) -> None: frames = [ AudioFrame(b"\xff\x7f", 16000, 1, 0, 1, {"duration_ms": 20, "speech": True, "speaker_id": "owner", "transcript": "你是谁"}), AudioFrame(b"\xff\x7f", 16000, 1, 20, 2, {"duration_ms": 20, "speech": True, "speaker_id": "owner"}), AudioFrame(b"\xff\x7f", 16000, 1, 40, 3, {"duration_ms": 20, "speech": True, "speaker_id": "background"}), AudioFrame(b"\xff\x7f", 16000, 1, 60, 4, {"duration_ms": 20, "speech": True, "speaker_id": "background"}), AudioFrame(b"\xff\x7f", 16000, 1, 80, 5, {"duration_ms": 20, "speech": True, "speaker_id": "background"}), ] stt = RecordingMetadataSttProvider() runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, memory_enabled=False, tool_router_enabled=False, noise_filter_enabled=False, endpoint_mode="primary_speaker", speaker_profile_ms=40, speaker_profile_min_ms=40, speaker_absent_ms=40, vad_min_duration_ms=1000, vad_end_silence_ms=1000, vad_max_recording_ms=200, end_chime_enabled=False, ), transport=MemoryAudioTransport(frames, flush_clears_input=False), stt=stt, llm=MockLlmProvider(["好的。"]), tts=SineTtsProvider(), ) summary = runtime.run(once=True) self.assertEqual(summary.completed_turns, 1) self.assertIsInstance(runtime.vad_recorder, PrimarySpeakerVadRecorder) self.assertEqual(stt.calls[0].metadata["end_reason"], "primary_speaker_absent") self.assertNotEqual(stt.calls[0].metadata["end_reason"], "max_recording") self.assertEqual(runtime.context.messages()[0].content, "你是谁") def test_run_agent_live_interrupts_playback_from_audio_hub_barge_in(self) -> None: initial_frames = [ AudioFrame(b"\xff\x7f", 16000, 1, 0, 1, {"duration_ms": 20, "speech": True, "speaker_id": "owner", "transcript": "介绍一下你自己"}), AudioFrame(b"\xff\x7f", 16000, 1, 20, 2, {"duration_ms": 20, "speech": True, "speaker_id": "owner", "transcript": "介绍一下你自己"}), AudioFrame(b"\x00\x00", 16000, 1, 40, 3, {"duration_ms": 20, "speech": False, "transcript": "介绍一下你自己"}), AudioFrame(b"\x00\x00", 16000, 1, 60, 4, {"duration_ms": 20, "speech": False, "transcript": "介绍一下你自己"}), ] interrupt_frames = [ AudioFrame(b"\xff\x7f", 16000, 1, 100, 10, {"duration_ms": 20, "speech": True, "speaker_id": "owner", "transcript": "等一下"}), AudioFrame(b"\xff\x7f", 16000, 1, 120, 11, {"duration_ms": 20, "speech": True, "speaker_id": "owner", "transcript": "等一下"}), ] transport = PlaybackInjectedTransport(initial_frames, injected_frames=interrupt_frames) runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, memory_enabled=False, tool_router_enabled=False, noise_filter_enabled=False, vad_min_duration_ms=40, vad_end_silence_ms=40, barge_in_enabled=True, barge_in_min_speech_ms=40, barge_in_echo_guard_ms=0, interrupt_target_latency_ms=100, end_chime_enabled=False, ), transport=transport, vad_recorder=VadRecorder(EnergyVadProvider(), min_duration_ms=40, end_silence_ms=40), stt=MetadataSttProvider(), llm=MockLlmProvider(["这是一个足够长的回答,用来验证播放时的后台麦克风打断能够停止当前播报。"]), tts=SineTtsProvider(), ) summary = runtime.run(once=True) self.assertTrue(summary.interrupted) self.assertTrue(runtime.cancellation_graph.root.cancelled) self.assertTrue(runtime._pending_capture_frames) def test_fake_apm_echo_does_not_trigger_interruption(self) -> None: fixture = build_fake_full_duplex_audio_fixture() apm = FakeWebRtcAudioProcessingProvider() detector = InterruptionDetector(vad=FakeVadProvider(), min_speech_ms=20) render = fixture[0] apm.process_render(render) processed_echo = apm.process_capture(fixture[0]) decision = detector.accept( processed_echo, state=PipelineState.SPEAKING, stt_events=[TranscriptEvent("stable_partial", "助手", is_stable=True)], ) self.assertTrue(processed_echo.metadata["echo_suppressed"]) self.assertFalse(decision.interrupted) def test_speaking_interruption_cancels_response_and_returns_to_listening(self) -> None: machine = FullDuplexStateMachine() graph = CancellationGraph("turn") detector = InterruptionDetector(vad=FakeVadProvider(), min_speech_ms=200) machine.transition(PipelineState.LISTENING, event_type="start") machine.transition(PipelineState.THINKING, event_type="final_transcript") machine.transition(PipelineState.SPEAKING, event_type="first_tts_chunk") detector.accept( build_fake_full_duplex_audio_fixture()[1], state=PipelineState.SPEAKING, stt_events=[TranscriptEvent("partial", "你", is_stable=False)], ) decision = detector.accept( build_fake_full_duplex_audio_fixture()[2], state=PipelineState.SPEAKING, stt_events=[TranscriptEvent("stable_partial", "你好", is_stable=True)], ) if decision.interrupted: graph.cancel_all("user interrupted") machine.transition(PipelineState.INTERRUPTED, event_type="interrupt_detected") machine.transition(PipelineState.LISTENING, event_type="buffered_user_audio") self.assertTrue(decision.interrupted) self.assertTrue(graph.root.cancelled) self.assertEqual(machine.current_state, PipelineState.LISTENING) def test_full_duplex_runtime_interrupt_fixture_uses_audio_hub_and_buffers_user_audio(self) -> None: runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, barge_in_min_speech_ms=200, ) ) fixture = build_fake_full_duplex_audio_fixture()[1:3] summary = runtime.run_interrupt_fixture(fixture, initial_state=PipelineState.SPEAKING) self.assertTrue(summary.interrupted) self.assertTrue(runtime.cancellation_graph.root.cancelled) self.assertEqual(runtime.state_machine.current_state, PipelineState.LISTENING) self.assertIsNotNone(runtime.interrupt_controller) self.assertEqual( [item.frame_id for item in runtime.interrupt_controller.buffered_user_frames], [2, 3], ) def test_streaming_stt_llm_tts_playback_order(self) -> None: stt = FakeStreamingSttProvider( scripted_events=[ [TranscriptEvent("partial", "你", is_stable=False)], [TranscriptEvent("stable_partial", "你好", is_stable=True)], ], final_text="你好", ) stt_session = stt.start_session("turn") llm = FakeStreamingLlmProvider([LlmStreamEvent("delta", "你好。"), LlmStreamEvent("finish", finish_reason="stop")]) segmenter = SentenceSegmenter() tts_session = FakeStreamingTtsProvider().start_stream(voice="default", sample_rate=16000) playback = InterruptiblePlaybackQueue() render = RenderReferenceRingBuffer(capacity_ms=1000) graph = CancellationGraph("turn") event_order: list[str] = [] for frame in build_fake_full_duplex_audio_fixture()[1:3]: for event in stt_session.accept_audio(frame): event_order.append(event.kind) final = stt_session.finish() event_order.append(final.kind) for llm_event in llm.stream([Message("user", final.text, 1.0)], cancellation=graph.root): event_order.append(f"llm_{llm_event.kind}") if llm_event.text_delta: for sentence in segmenter.accept_delta(llm_event.text_delta): event_order.append("sentence") frames = tts_session.accept_text(sentence) event_order.append("tts") playback.enqueue(sentence, frames) result = playback.play_next(render_reference=render, cancellation=graph.root) event_order.append("playback") self.assertEqual( event_order, ["partial", "stable_partial", "final", "llm_delta", "sentence", "tts", "llm_finish", "playback"], ) self.assertFalse(result.interrupted) self.assertEqual(playback.spoken.text, "你好。") self.assertEqual(render.frame_count, 1) def test_full_duplex_runtime_streaming_response_writes_render_reference_and_spoken_text(self) -> None: runtime = FullDuplexAgentRuntime( config=AppConfig(audio_apm_provider="fake", audio_apm_required=False), llm_provider=FakeStreamingLlmProvider( [LlmStreamEvent("delta", "第一句。第二句。"), LlmStreamEvent("finish", finish_reason="stop")] ), tts_provider=FakeStreamingTtsProvider(), ) spoken = runtime.run_streaming_response_fixture([Message("user", "你好", 1.0)]) self.assertEqual(spoken, "第一句。第二句。") self.assertIsNotNone(runtime.audio_hub) self.assertEqual(runtime.audio_hub.render_reference.frame_count, 2) self.assertEqual(runtime.playback_queue.pending_items, 0) def test_full_duplex_runtime_injects_memory_context_before_llm(self) -> None: memory = FakeMemoryManager() memory.save(MemoryRecordInput("preference", "用户喜欢 Python")) llm = FakeStreamingLlmProvider([LlmStreamEvent("delta", "记住了。"), LlmStreamEvent("finish", finish_reason="stop")]) runtime = FullDuplexAgentRuntime( config=AppConfig(audio_apm_provider="fake", audio_apm_required=False, memory_enabled=True), memory_manager=memory, llm_provider=llm, tts_provider=FakeStreamingTtsProvider(), ) spoken = runtime.run_conversation_response_fixture("Python 项目怎么做?") self.assertEqual(spoken, "记住了。") self.assertIn("长期记忆", llm.requests[0][1].content) self.assertIn("用户喜欢 Python", llm.requests[0][1].content) self.assertEqual([message.role for message in runtime.context.messages()], ["user", "assistant"]) def test_full_duplex_runtime_memory_health_detects_manifest_mismatch(self) -> None: with tempfile.TemporaryDirectory() as tmp: memory = SQLiteMemoryManager(Path(tmp) / "memory.sqlite3") saved = memory.save(MemoryRecordInput("fact", "Owner 正在做全双工语音助手")) broken = FaissIndexManifest( embedding_model="fake", record_ids=(saved.id, "missing-id"), checksums={saved.id: "wrong"}, ) runtime = FullDuplexAgentRuntime( config=AppConfig(audio_apm_provider="fake", audio_apm_required=False, memory_enabled=True), memory_manager=memory, memory_manifest=broken, ) health = runtime.check_memory_health() self.assertFalse(health.ok) self.assertTrue(any("missing-id" in error for error in health.errors)) self.assertTrue(any("checksum mismatch" in error for error in health.errors)) def test_full_duplex_runtime_routes_memory_search_tool_call(self) -> None: memory = FakeMemoryManager() memory.save(MemoryRecordInput("project", "Owner 项目正在做全双工 Agent")) runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, memory_enabled=True, tool_router_enabled=True, ), memory_manager=memory, llm_provider=FakeStreamingLlmProvider( [ LlmStreamEvent( "tool_call", tool_call={ "id": "tool-1", "name": "memory.search", "arguments": {"query": "Owner Agent", "top_k": 1}, "turn_id": "turn-1", }, ), LlmStreamEvent("delta", "查到了。"), LlmStreamEvent("finish", finish_reason="stop"), ] ), tts_provider=FakeStreamingTtsProvider(), ) spoken = runtime.run_conversation_response_fixture("查一下当前项目") self.assertEqual(spoken, "查到了。") self.assertEqual(runtime.tool_results[0].status, "success") self.assertIn("全双工 Agent", runtime.tool_results[0].output_text) self.assertIn("工具结果 memory.search", runtime.tool_result_messages[0].content) def test_full_duplex_runtime_high_risk_tool_call_requires_confirmation(self) -> None: runtime = FullDuplexAgentRuntime( config=AppConfig( audio_apm_provider="fake", audio_apm_required=False, memory_enabled=True, tool_router_enabled=True, ), memory_manager=FakeMemoryManager(), llm_provider=FakeStreamingLlmProvider( [ LlmStreamEvent( "tool_call", tool_call={ "id": "tool-1", "name": "memory.search", "arguments": {"query": "账号"}, "natural_language_intent": "上传账号资料", "turn_id": "turn-1", }, ), LlmStreamEvent("finish", finish_reason="stop"), ] ), tts_provider=FakeStreamingTtsProvider(), ) runtime.run_conversation_response_fixture("上传账号资料") self.assertEqual(runtime.tool_results[0].status, "confirmation_required") self.assertEqual(runtime.tool_router.audit_log[0].action, "require_confirmation") def test_memory_restart_and_tool_search_integration(self) -> None: with tempfile.TemporaryDirectory() as tmp: db_path = Path(tmp) / "memory.sqlite3" SQLiteMemoryManager(db_path).save(MemoryRecordInput("preference", "用户喜欢 Python")) restarted = SQLiteMemoryManager(db_path) router = ToolRouter({"memory.search": MemorySearchTool()}) request = ToolCallRequest("1", "memory.search", {"query": "Python"}, "turn") decision = router.route(request, ToolContext(memory=restarted)) result = router.execute(request, decision, ToolContext(memory=restarted)) self.assertEqual(result.status, "success") self.assertIn("用户喜欢 Python", result.output_text) def test_tool_router_security_blocks_high_risk_fake_integration(self) -> None: router = ToolRouter({"memory.search": MemorySearchTool()}) request = ToolCallRequest( "1", "memory.search", {"query": "账号"}, "turn", natural_language_intent="上传账号资料", ) decision = router.route(request, ToolContext(memory=FakeMemoryManager())) self.assertEqual(decision.action, "require_confirmation") self.assertEqual(decision.risk_level, "high") def test_performance_metrics_are_sanitized(self) -> None: recorder = PerformanceMetricRecorder() recorder.record( "interrupt_latency", started_at_ms=100, finished_at_ms=250, payload={"api_key": "secret", "preview": "tp-" + "abcdefghijklmnop"}, ) self.assertEqual(recorder.summary()["interrupt_latency"], 150) self.assertEqual(recorder.metrics[0].payload["api_key"], "[redacted]") self.assertFalse(diagnostics_contain_sensitive_data(recorder.metrics[0].payload)) def test_sanitize_diagnostics_removes_nested_sensitive_values(self) -> None: sanitized = sanitize_diagnostics( { "nested": {"authorization": "Bearer secret", "raw_audio": b"bytes"}, "text": "normal", } ) self.assertEqual(sanitized["nested"]["authorization"], "[redacted]") self.assertEqual(sanitized["nested"]["raw_audio"], "[redacted]") self.assertFalse(diagnostics_contain_sensitive_data(sanitized)) if __name__ == "__main__": unittest.main()