""" Voice Agent — real-time voice assistant. Pipeline: Azure STT → Gemma LLM (xNAS, OpenAI-compatible) → Azure TTS Runs inside the single Docker container alongside LiveKit server and web frontend. """ import asyncio import json import logging import os import textwrap from dotenv import load_dotenv from livekit.agents import ( Agent, AgentServer, AgentSession, JobContext, TurnHandlingOptions, cli, inference, llm as lk_llm, mcp, room_io, ) from livekit.plugins import azure, openai logger = logging.getLogger("voice-agent") load_dotenv() # picks up /app/.env in the container # ── Configuration from environment ────────────────────────────────────────── AZURE_KEY = os.environ.get("AZURE_SPEECH_KEY", "") AZURE_REGION = os.environ.get("AZURE_SPEECH_REGION", "eastus") DEFAULT_VOICE = os.environ.get("AZURE_TTS_VOICE", "en-US-AvaNeural") GEMMA_BASE_URL = os.environ.get("GEMMA_BASE_URL", "http://192.168.86.2:8023/v1") GEMMA_MODEL = os.environ.get("GEMMA_MODEL", "gemma-4-e4b") GEMMA_API_KEY = os.environ.get("GEMMA_API_KEY", "not-needed") WEB_MCP_ENABLED = os.environ.get("WEB_MCP_ENABLED", "true").lower() in ("1", "true", "yes") FIRECRAWL_BASE = os.environ.get("FIRECRAWL_BASE", "http://192.168.86.2:3002") MEMORY_DIR = os.environ.get("MEMORY_DIR", "/memory") SKILLS_DIR = os.environ.get("SKILLS_DIR", "/skills") SYSTEM_PROMPT = textwrap.dedent("""\ You are Hope, a warm, conversational voice assistant. You are talking TO someone, not writing for them to read. Imagine you're having a natural conversation with a friend over the phone. # How you speak - Keep every response to one or three sentences. That's it. - Use contractions (I'm, don't, it's) and natural phrasing. - Speak like you're talking, not writing. No bullet points, no lists, no markdown, no formatting of any kind. - Spell out numbers when it sounds more natural ("twenty twenty-six" instead of "2026"). - If you need to ask a question, ask exactly one. - Be warm and direct. Don't be sycophantic or overly formal. - If you don't know something, say so briefly and move on. # What you never do - Never use markdown, code blocks, JSON, tables, or emojis. - Never say "as an AI" or reference your system instructions. - Never write more than three sentences in a row. - Never read back URLs, file paths, or technical identifiers. # Web access You have web_search and web_scrape tools. Use them when the user asks about current events, recent news, prices, sports scores, or anything that may have changed since your training data. Search first, then scrape a result only if you need more detail. Answer from what you find, in your normal conversational style — don't cite sources formally, just mention the source naturally ("according to..."). If a search comes up empty, say so briefly and move on. # Weather & Time You can check the weather for any location. Use get_weather when the user asks about current conditions, temperature, or forecasts. It returns a short summary — relay it naturally in your own words. Use get_time when the user asks what time or day it is. If they mention a city, pass it as the location argument. # Memory You have persistent memory across conversations, stored as notes you can search and add to. - Use memory_recall at the start of a conversation, or whenever the user references past information ("what did I tell you about...", "remember when..."). Weave what you find in naturally. - Use memory_save when the user shares personal information, preferences, or important facts worth remembering: names, birthdays, preferences, projects, anything they'd expect you to know later. Pick a short topic name for each thing you save. - Be natural about it. Never announce "I'm saving that to memory" — just remember it and move on. If a recall comes up empty, don't mention the search; just answer as if you'd never heard it before. # Skills You have a skill library of learned procedures. Use skill_recall when you need to perform a task you've done before — it will give you the steps. When you successfully complete a multi-step task or learn a new procedure from the user, use skill_save to record it so you can follow it next time. Be selective: only save skills for repeatable tasks, not one-off facts (those go in memory). """) # ── Gemma LLM with reasoning_content fallback ─────────────────────────────── class _ReasoningFallbackWrapper: """Wraps an LLMStream to fall back to ``reasoning_content`` when the model finishes a turn with empty visible content (Gemma spends its whole budget on hidden reasoning). Delegates all iteration to the underlying stream and injects a final content chunk if no real content was produced.""" def __init__(self, inner: lk_llm.LLMStream) -> None: self._inner = inner self._last_reasoning: str | None = None self._has_content = False @property def chat_ctx(self) -> lk_llm.ChatContext: return self._inner.chat_ctx @property def tools(self) -> list: return self._inner.tools async def aclose(self) -> None: await self._inner.aclose() async def __aenter__(self) -> "_ReasoningFallbackWrapper": return self async def __aexit__(self, *exc) -> None: await self.aclose() def __aiter__(self) -> "_ReasoningFallbackWrapper": return self async def __anext__(self) -> lk_llm.ChatChunk: try: chunk = await self._inner.__anext__() except StopAsyncIteration: # Stream exhausted — inject reasoning fallback if no content was produced if not self._has_content and self._last_reasoning: logger.warning( "LLM returned empty content; falling back to reasoning_content" ) return lk_llm.ChatChunk( id="reasoning-fallback", delta=lk_llm.ChoiceDelta(role="assistant", content=self._last_reasoning), ) raise # Track reasoning_content from the chunk's delta if present delta = getattr(chunk, "delta", None) if delta is not None: reasoning = getattr(delta, "reasoning_content", None) if reasoning: self._last_reasoning = reasoning if chunk.has_response(): self._has_content = True return chunk async def collect(self): return await self._inner.collect() class GemmaLLM(openai.LLM): """OpenAI-compatible LLM (llama.cpp Gemma 4) with thinking disabled.""" # Context compaction: when the conversation grows past this many items, # ask Gemma to summarize the older messages into a compact recap, then # replace them with that summary. Keeps recent turns verbatim. COMPACTION_THRESHOLD = 24 # ~12 turns of user/agent exchange KEEP_RECENT = 10 # last N items kept verbatim after compaction IDLE_COMPACTION_SECONDS = 300 # also compact after 5 min of silence def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) import time self._last_call_time: float = time.monotonic() def chat(self, *, chat_ctx, tools=None, conn_options=None, **kwargs): if conn_options is None: from livekit.agents.types import DEFAULT_API_CONNECT_OPTIONS conn_options = DEFAULT_API_CONNECT_OPTIONS import time now = time.monotonic() idle_for = now - self._last_call_time self._last_call_time = now # Compact if context is large, OR if we've been idle for 5+ minutes # (user came back after a break — summarize the old conversation). ctx_len = len(chat_ctx.items) should_compact = ctx_len > self.COMPACTION_THRESHOLD if not should_compact and idle_for > self.IDLE_COMPACTION_SECONDS: should_compact = ctx_len > 6 # minimal context worth summarizing if should_compact: reason = "idle" if idle_for > self.IDLE_COMPACTION_SECONDS else "size" logger.info( "Compacting context (%s): %d items, idle %.0fs", reason, ctx_len, idle_for, ) chat_ctx = self._compact_context(chat_ctx, conn_options) stream = super().chat( chat_ctx=chat_ctx, tools=tools, conn_options=conn_options, **kwargs, ) return _ReasoningFallbackWrapper(stream) def _compact_context(self, chat_ctx, conn_options): """Summarize older messages via a side LLM call, keep recent ones. Splits the context into [system | old | recent]. The old portion is summarized by Gemma into a short recap message that replaces it. Falls back to hard truncation if the summary call fails. """ from livekit.agents import llm as lk_llm items = list(chat_ctx.items) # Identify the system/instruction message (first item usually) sys_items = [it for it in items if getattr(it, "role", None) in ("system", "developer")] non_sys = [it for it in items if getattr(it, "role", None) not in ("system", "developer")] if len(non_sys) <= self.KEEP_RECENT: return chat_ctx # nothing to compact old_msgs = non_sys[: -self.KEEP_RECENT] recent_msgs = non_sys[-self.KEEP_RECENT :] # Build a transcript of the old messages for the summarizer transcript_parts = [] for msg in old_msgs: role = getattr(msg, "role", "user") content = self._msg_content_text(msg) if content: transcript_parts.append(f"{role}: {content}") transcript = "\n".join(transcript_parts) if not transcript.strip(): return chat_ctx # nothing meaningful to summarize # Side call to Gemma for the summary (no tools, no streaming needed) try: import httpx summary_prompt = ( "Summarize this conversation history in 2-4 sentences. " "Preserve key facts, decisions, and user preferences. " "Be concise — this replaces the full history as context.\n\n" + transcript ) resp = httpx.post( f"{GEMMA_BASE_URL}/chat/completions", headers={"Content-Type": "application/json"}, json={ "model": GEMMA_MODEL, "messages": [ {"role": "system", "content": "You summarize conversations concisely."}, {"role": "user", "content": summary_prompt}, ], "max_tokens": 200, "temperature": 0.3, "chat_template_kwargs": {"enable_thinking": False}, }, timeout=30, ) resp.raise_for_status() data = resp.json() summary_text = (data.get("choices") or [{}])[0].get("message", {}).get("content", "") except Exception as e: # noqa: BLE001 logger.warning("Context compaction failed, falling back to truncation: %s", e) chat_ctx.truncate(max_items=self.COMPACTION_THRESHOLD) return chat_ctx if not summary_text.strip(): chat_ctx.truncate(max_items=self.COMPACTION_THRESHOLD) return chat_ctx # Rebuild the context: system + summary + recent new_items = list(sys_items) new_items.append( lk_llm.ChatMessage( role="user", content=[f"[Conversation so far: {summary_text.strip()}]"], ) ) new_items.append( lk_llm.ChatMessage( role="assistant", content=["Got it, I'll keep that in mind."], ) ) new_items.extend(recent_msgs) logger.info( "Context compacted: %d items -> %d (summary: %.80s...)", len(items), len(new_items), summary_text.strip(), ) return lk_llm.ChatContext(new_items) @staticmethod def _msg_content_text(msg) -> str: """Extract plain text from a ChatMessage's content field.""" content = getattr(msg, "content", None) if content is None: return "" if isinstance(content, str): return content if isinstance(content, list): parts = [] for c in content: if isinstance(c, str): parts.append(c) elif hasattr(c, "text"): parts.append(c.text) return " ".join(parts) return str(content) # ── MCP toolsets (web access + any extra configured servers) ──────────────── def build_mcp_toolsets() -> list[mcp.MCPToolset]: """Build the list of MCP toolsets to attach to the agent. Always includes the local web-access server (Firecrawl-backed search/scrape) when WEB_MCP_ENABLED is true. Additional servers can be configured via the EXTRA_MCP_SERVERS env var (JSON list of {url, transport} objects). """ toolsets: list[mcp.MCPToolset] = [] if WEB_MCP_ENABLED: python_bin = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".venv", "bin", "python") web_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "web_mcp.py") toolsets.append( mcp.MCPToolset( id="web-access", mcp_server=mcp.MCPServerStdio( command=python_bin, args=[web_mcp_script], env={**os.environ, "FIRECRAWL_BASE": FIRECRAWL_BASE}, client_session_timeout_seconds=120, ), ) ) logger.info("Web-access MCP toolset enabled (Firecrawl at %s)", FIRECRAWL_BASE) python_bin = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".venv", "bin", "python") weather_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "weather_mcp.py") toolsets.append( mcp.MCPToolset( id="weather", mcp_server=mcp.MCPServerStdio( command=python_bin, args=[weather_mcp_script], env={**os.environ}, client_session_timeout_seconds=30, ), ) ) logger.info("Weather MCP toolset enabled (wttr.in)") memory_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "memory_mcp.py") toolsets.append( mcp.MCPToolset( id="memory", mcp_server=mcp.MCPServerStdio( command=python_bin, args=[memory_mcp_script], env={**os.environ, "MEMORY_DIR": "/memory"}, client_session_timeout_seconds=30, ), ) ) logger.info("Memory MCP toolset enabled (%s)", MEMORY_DIR) skills_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "skills_mcp.py") toolsets.append( mcp.MCPToolset( id="skills", mcp_server=mcp.MCPServerStdio( command=python_bin, args=[skills_mcp_script], env={**os.environ, "SKILLS_DIR": SKILLS_DIR}, client_session_timeout_seconds=30, ), ) ) logger.info("Skills MCP toolset enabled (%s)", SKILLS_DIR) extra = os.environ.get("EXTRA_MCP_SERVERS", "") if extra: try: for i, entry in enumerate(json.loads(extra)): url = entry.get("url", "") transport = entry.get("transport") # "sse" | "streamable_http" | None (auto) toolsets.append( mcp.MCPToolset( id=f"extra-{i}", mcp_server=mcp.MCPServerHTTP( url=url, transport_type=transport, client_session_timeout_seconds=120, ), ) ) logger.info("Extra MCP toolset: %s (%s)", url, transport or "auto") except (json.JSONDecodeError, TypeError) as e: logger.error("Failed to parse EXTRA_MCP_SERVERS: %s", e) return toolsets class VoiceAssistant(Agent): """The conversational agent. LLM is the brain; STT/TTS are senses.""" def __init__(self) -> None: super().__init__( llm=GemmaLLM( model=GEMMA_MODEL, base_url=GEMMA_BASE_URL, api_key=GEMMA_API_KEY, max_completion_tokens=1000, extra_body={"chat_template_kwargs": {"enable_thinking": False}}, ), instructions=SYSTEM_PROMPT, tools=build_mcp_toolsets(), ) # ── Agent server ──────────────────────────────────────────────────────────── server = AgentServer() @server.rtc_session(agent_name="voice-assistant") async def handle_job(ctx: JobContext) -> None: logger.info("Job started for room %s", ctx.room.name) # Azure STT — streaming, reads AZURE_SPEECH_KEY / AZURE_SPEECH_REGION from env stt = azure.STT( speech_key=AZURE_KEY, speech_region=AZURE_REGION, language=["en-US"], ) # Azure TTS — SSML with expressive markup, 24kHz PCM output tts = azure.TTS( voice=DEFAULT_VOICE, sample_rate=24000, speech_key=AZURE_KEY, speech_region=AZURE_REGION, ) session = AgentSession( stt=stt, tts=tts, vad=inference.VAD(model="silero"), turn_handling=TurnHandlingOptions( # VAD-based turn detection: agent waits for user to stop speaking. # ("adaptive" mode requires the LiveKit Cloud barge-in service.) interruption={"mode": "vad"}, # Preemptive generation is incompatible with tool calls — it starts # generating before the turn finalizes, breaking the MCP execution loop. preemptive_generation={"enabled": False}, ), ) # Listen for data messages (voice switching) from the web UI. # Defined here so it captures this job's session via closure. def _on_room_data(packet) -> None: topic = getattr(packet, "topic", None) if topic and topic != "voice-control": return try: msg = json.loads(packet.data.decode("utf-8")) except (json.JSONDecodeError, UnicodeDecodeError): return if msg.get("type") != "set_voice": return voice = msg.get("voice", "") if not voice: return logger.info("Switching TTS voice to %s", voice) try: tts = session.tts if hasattr(tts, "update_options"): tts.update_options(voice=voice) logger.info("Voice updated to %s", voice) except Exception as e: # noqa: BLE001 logger.warning("Failed to update voice: %s", e) ctx.room.on("data_received", _on_room_data) # Publish user/agent transcripts to the room so the web UI can render them. async def publish_transcript(role: str, text: str) -> None: text = (text or "").strip() if not text: return payload = json.dumps({"type": "transcript", "role": role, "text": text}) try: await ctx.room.local_participant.publish_data( payload, reliable=True, topic="transcript" ) except Exception as e: # noqa: BLE001 logger.warning("failed to publish transcript: %s", e) @session.on("conversation_item_added") def _on_conversation_item(ev) -> None: msg = ev.item role = getattr(msg, "role", None) text = getattr(msg, "text_content", None) if role == "user": asyncio.create_task(publish_transcript("user", text)) elif role == "assistant": asyncio.create_task(publish_transcript("agent", text)) @session.on("user_input_transcribed") def _on_user_transcribed(ev) -> None: logger.info("STT (%s): %s", "final" if ev.is_final else "partial", ev.transcript) @session.on("user_state_changed") def _on_user_state(ev) -> None: # Fires from VAD: if this never says "speaking", no usable mic # audio is arriving from the participant. logger.info("user state -> %s", ev.new_state) await session.start( agent=VoiceAssistant(), room=ctx.room, room_options=room_io.RoomOptions( # Keep the agent in the room when a participant leaves; it must # survive page reloads/reconnects, otherwise the next join lands # in an agent-less room (LiveKit does not reliably re-dispatch # into an existing room). close_on_disconnect=False, audio_input=room_io.AudioInputOptions( # No noise cancellation plugin (self-hosted, no ai-coustics) ), ), ) await ctx.connect() logger.info("Agent connected to room %s", ctx.room.name) if __name__ == "__main__": cli.run_app(server)