"""Shared client-visible output for deterministic fixed speech.""" from __future__ import annotations from collections.abc import Awaitable from typing import Any from pipecat.frames.frames import OutputTransportMessageUrgentFrame, TTSSpeakFrame from pipecat.utils.time import time_now_iso8601 from services.brains.base import BrainRuntime from services.runtime_variables import DynamicVariableStore class FixedSpeechOutput: """Display and synthesize fixed speech without waiting for playback.""" def __init__( self, store: DynamicVariableStore, runtime: BrainRuntime, ) -> None: self._store = store self._runtime = runtime self._client_ready = False self._pending_transcripts: list[dict[str, Any]] = [] async def mark_client_ready(self) -> None: self._client_ready = True pending = self._pending_transcripts self._pending_transcripts = [] for message in pending: await self.emit(message) async def speak( self, text: str, *, source: str, node_id: str | None = None, record_history: bool = True, ) -> Awaitable[None] | None: content = text.strip() if not content: return None if record_history: self._store.record("agent", content) transcript = { "type": "transcript", "role": "assistant", "content": content, "timestamp": time_now_iso8601(), "source": source, **({"nodeId": node_id} if node_id else {}), } if self._client_ready: await self.emit(transcript) else: self._pending_transcripts.append(transcript) track_speech = getattr(self._runtime.call_end, "track_speech", None) playback_completion: Awaitable[None] | None = None if callable(track_speech): playback_completion = track_speech() await self._runtime.queue_frame( TTSSpeakFrame(content, append_to_context=False) ) return playback_completion async def emit(self, message: dict[str, Any]) -> None: await self._runtime.queue_frame( OutputTransportMessageUrgentFrame(message=message) )