471 lines
22 KiB
Python
471 lines
22 KiB
Python
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()
|