"""Praxis Pipecat server pipeline — minimal viable voice loop (SLICE-02 TASK-02-04). Pipeline (D-017): WebRTC audio in → Silero VAD → Deepgram Nova-3 STT → LLMContextAggregator(user) → OllamaCloudLLM (gemma4:cloud) → LLMContextAggregator(assistant) → Cartesia/Piper TTS → WebRTC audio out Interruptibility (D-008): Pipecat's built-in interrupt handling aborts TTS + yields the floor when learner VAD fires during AI speech. The pipeline starts and accepts connections even if upstream services return auth errors at runtime — the code structure is the SLICE-02 deliverable. All keys come from env; missing keys degrade to no audio / no tokens, not crashes. Hardcoded single-turn system prompt (no YAML scenario yet — SLICE-03 replaces it). """ from __future__ import annotations import os from typing import Any from loguru import logger def _env(key: str, default: str = "") -> str: return os.environ.get(key, default).strip() # Hardcoded single-turn system prompt (SLICE-02 walking skeleton). # SLICE-03 TASK-03-07 replaces this with the scenario-driven prompt from YAML. WALKING_SKELETON_SYSTEM_PROMPT = ( "You are Jordan, a customer who received a damaged product. " "You are frustrated but not abusive. You want a refund. " "Stay in character. Do not break role. " "Keep responses concise for voice (1-3 sentences)." ) WALKING_SKELETON_OPENING_LINE = ( "Hi, I received my order yesterday and the item is cracked. I want my money back." ) def _build_llm_context(scenario_runtime=None): """Build the LLMContext with the scenario-driven system prompt (TASK-03-07). If a scenario_runtime is provided, uses scenario.setup.system_prompt. Otherwise falls back to the SLICE-02 walking-skeleton prompt. """ from pipecat.processors.aggregators.llm_context import LLMContext if scenario_runtime is not None: system_prompt = scenario_runtime.system_prompt else: system_prompt = WALKING_SKELETON_SYSTEM_PROMPT messages = [ {"role": "system", "content": system_prompt}, ] return LLMContext(messages=messages) def _build_transport(webrtc_connection) -> Any: """Build the SmallWebRTCTransport with audio in/out enabled.""" from pipecat.transports.base_transport import TransportParams from pipecat.transports.smallwebrtc.transport import SmallWebRTCTransport params = TransportParams( audio_in_enabled=True, audio_out_enabled=True, audio_out_sample_rate=24000, ) return SmallWebRTCTransport(webrtc_connection, params) def _build_stt() -> Any: """Build the Deepgram Nova-3 STT service (D-013).""" from pipecat.services.deepgram.stt import DeepgramSTTService api_key = _env("DEEPGRAM_API_KEY") if not api_key: logger.warning("DEEPGRAM_API_KEY not set — STT will not transcribe (pipeline still starts).") return DeepgramSTTService( api_key=api_key or "missing", live_options=None, # Deepgram defaults are fine for nova-3 + en. ) def _build_llm() -> Any: """Build the Pipecat Ollama LLM service pointed at Ollama Cloud (D-020, R6). Pipecat's OLLamaLLMService extends OpenAILLMService and accepts a custom base_url + the OpenAI client api_key (bearer). We point it at https://ollama.com/v1 with OLLAMA_API_KEY as the bearer. """ from pipecat.services.ollama.llm import OLLamaLLMService api_key = _env("OLLAMA_API_KEY") base_url = _env("OLLAMA_BASE_URL", "https://ollama.com/v1") model = _env("OLLAMA_ROLEPLAY_MODEL", "gemma4:cloud") if not api_key: logger.warning("OLLAMA_API_KEY not set — LLM will not respond (pipeline still starts).") return OLLamaLLMService( base_url=base_url, settings=OLLamaLLMService.Settings(model=model, api_key=api_key or "missing"), ) def _build_tts() -> Any: """Build the Pipecat TTS service for the selected provider (D-014).""" choice = _env("PRAXIS_TTS", "cartesia").lower() if choice == "piper": from pipecat.services.piper.tts import PiperTTSService voice_model = _env("PIPER_VOICE_MODEL") if not voice_model: logger.warning("PIPER_VOICE_MODEL not set — Piper TTS will not speak (pipeline still starts).") return PiperTTSService( voice_id=voice_model or "missing", ) # Default: Cartesia from pipecat.services.cartesia.tts import CartesiaTTSService api_key = _env("CARTESIA_API_KEY") voice_id = _env("CARTESIA_VOICE_ID", "a3536a36-1d18-4efb-a95a-7c44b7b5e384") if not api_key: logger.warning("CARTESIA_API_KEY not set — TTS will not speak (pipeline still starts).") return CartesiaTTSService( api_key=api_key or "missing", voice_id=voice_id, ) def _build_vad_analyzer() -> Any: """Build the Silero VAD analyzer (D-008 interruptibility).""" from pipecat.audio.vad.silero import SileroVADAnalyzer return SileroVADAnalyzer() def build_pipeline(webrtc_connection, *, scenario_id: str | None = None): """Assemble the full Pipecat pipeline + task + runner for one WebRTC session. Args: webrtc_connection: a SmallWebRTCConnection with an accepted offer. scenario_id: if set, load the scenario and use its system prompt + opening line (TASK-03-07). If None, falls back to the walking-skeleton prompt. Returns (pipeline, task, runner, transport, scenario_runtime) so the caller can start the task on connection, play the opening line, and run the branch classifier + debrief at session end. """ from pipecat.pipeline.pipeline import Pipeline from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.processors.aggregators.llm_response_universal import ( LLMContextAggregator, ) # Load the scenario runtime (TASK-03-03, TASK-03-07). scenario_runtime = None if scenario_id: try: from server.scenarios.runtime import build_runtime_from_id scenario_runtime = build_runtime_from_id(scenario_id) logger.info( f"Loaded scenario {scenario_id!r}: branches={scenario_runtime.scenario.branch_ids()}" ) except Exception as exc: logger.warning( f"Could not load scenario {scenario_id!r}: {exc}. " f"Falling back to walking-skeleton prompt." ) transport = _build_transport(webrtc_connection) stt = _build_stt() llm = _build_llm() tts = _build_tts() from server.latency import LatencyObserver latency_observer = LatencyObserver() context = _build_llm_context(scenario_runtime) user_aggregator = LLMContextAggregator(context=context, role="user") assistant_aggregator = LLMContextAggregator(context=context, role="assistant") pipeline = Pipeline( [ transport.input(), # WebRTC audio in stt, # Deepgram Nova-3 latency_observer, # timestamp ASR-ready (TASK-02-06) user_aggregator, # collect user transcript into context llm, # Ollama gemma4:cloud latency_observer, # timestamp LLM-first-token (passes through) tts, # Cartesia/Piper latency_observer, # timestamp TTS-first-audio + emit metric transport.output(), # WebRTC audio out assistant_aggregator, # collect assistant text into context ] ) task = PipelineTask( pipeline, params=PipelineParams( allow_interruptions=True, # D-008 abort-and-yield enable_metrics=True, # latency measurement (TASK-02-06) metrics_request_timeout=10.0, ), ) runner = PipelineRunner(handle_sigint=False) return pipeline, task, runner, transport, scenario_runtime def build_runtime_from_id(scenario_id: str): """Re-export of the scenario runtime builder (TASK-03-07).""" from server.scenarios.runtime import build_runtime_from_id as _br return _br(scenario_id) __all__ = [ "build_pipeline", "build_runtime_from_id", "WALKING_SKELETON_SYSTEM_PROMPT", "WALKING_SKELETON_OPENING_LINE", ]