diff --git a/docs/assistant_server_contract.md b/docs/assistant_server_contract.md index 62124ac5..91b23df8 100644 --- a/docs/assistant_server_contract.md +++ b/docs/assistant_server_contract.md @@ -115,7 +115,7 @@ The user simulator always connects to `ws://localhost:{port}/ws`, so both `/ws` The user simulator sends and receives audio in the [Twilio Media Streams](https://www.twilio.com/docs/voice/media-streams) format: JSON -envelopes over WebSocket. Helper functions in `audio_bridge.py` handle all encoding +envelopes over WebSocket. Helper functions in `utils/audio_utils.py` handle all encoding and decoding. ### Incoming messages from the simulator @@ -136,7 +136,7 @@ events. ### Sending audio back to the simulator ```python -from eva.assistant.audio_bridge import create_twilio_media_message +from eva.utils.audio_utils import create_twilio_media_message msg = create_twilio_media_message(stream_sid, mulaw_chunk) await websocket.send_text(msg) ``` @@ -147,7 +147,7 @@ dedicated `_pace_audio_output` asyncio task to drain a queue at this rate — co pattern. If audio is sent too fast or too slow, the user simulator's may incorrectly detect turn boundaries. -### Audio conversion utilities (`audio_bridge.py`) +### Audio conversion utilities (`utils/audio_utils.py`) | Function | Converts | |---|---| @@ -172,7 +172,7 @@ resulting mixed WAV will have the audio offset by Δ. Before extending either bu pad the *other* buffer to the same position: ```python -from eva.assistant.audio_bridge import sync_buffer_to_position +from eva.utils.audio_utils import sync_buffer_to_position # When user audio arrives and the model is not speaking: if not model_is_speaking: @@ -381,17 +381,16 @@ from pathlib import Path import uvicorn from fastapi import FastAPI, WebSocket -from eva.assistant.audio_bridge import ( - FrameworkLogWriter, - MetricsLogWriter, +from eva.assistant.base_server import INITIAL_MESSAGE, AbstractAssistantServer +from eva.assistant.pipeline.observers import FrameworkLogWriter, MetricsLogWriter +from eva.models.config import ModelConfig +from eva.utils.audio_utils import ( create_twilio_media_message, mulaw_8k_to_pcm16_24k, parse_twilio_media_message, pcm16_24k_to_mulaw_8k, sync_buffer_to_position, ) -from eva.assistant.base_server import INITIAL_MESSAGE, AbstractAssistantServer -from eva.models.config import ModelConfig class MyFrameworkAssistantServer(AbstractAssistantServer): diff --git a/src/eva/assistant/audio_bridge.py b/src/eva/assistant/audio_bridge.py deleted file mode 100644 index af8fc9af..00000000 --- a/src/eva/assistant/audio_bridge.py +++ /dev/null @@ -1,261 +0,0 @@ -"""Shared audio bridge utilities for framework-specific assistant servers. - -All framework servers need to: -1. Accept Twilio-framed WebSocket connections from the user simulator -2. Convert audio between Twilio's mulaw 8kHz and the framework's native format -3. Write framework_logs.jsonl with timestamped events - -This module provides the common infrastructure. -""" - -import audioop -import base64 -import json -import struct -import time -from pathlib import Path - -import numpy as np -import soxr - -from eva.utils.logging import get_logger - -logger = get_logger(__name__) - - -# ── Audio format conversion ────────────────────────────────────────── - - -def mulaw_8k_to_pcm16_16k(mulaw_bytes: bytes) -> bytes: - """Convert 8kHz mu-law audio to 16kHz 16-bit PCM.""" - # Decode mu-law to 16-bit PCM at 8kHz - pcm_8k = audioop.ulaw2lin(mulaw_bytes, 2) - # Upsample from 8kHz to 16kHz - pcm_16k, _ = audioop.ratecv(pcm_8k, 2, 1, 8000, 16000, None) - return pcm_16k - - -def mulaw_8k_to_pcm16_24k(mulaw_bytes: bytes) -> bytes: - """Convert 8kHz mu-law audio to 24kHz 16-bit PCM.""" - # Decode mu-law to 16-bit PCM at 8kHz - pcm_8k = audioop.ulaw2lin(mulaw_bytes, 2) - # Upsample from 8kHz to 24kHz - pcm_24k, _ = audioop.ratecv(pcm_8k, 2, 1, 8000, 24000, None) - # audioop.ratecv can produce ±2 samples; clamp to exact 3× input length - # so that the inverse conversion recovers the original sample count. - expected_bytes = len(pcm_8k) * 3 - if len(pcm_24k) < expected_bytes: - pcm_24k = pcm_24k + b"\x00" * (expected_bytes - len(pcm_24k)) - elif len(pcm_24k) > expected_bytes: - pcm_24k = pcm_24k[:expected_bytes] - return pcm_24k - - -def pcm16_24k_to_mulaw_8k(pcm_bytes: bytes) -> bytes: - """Convert 24kHz 16-bit PCM to 8kHz mu-law. - - Uses soxr VHQ resampling (same as Pipecat) for proper anti-aliasing during the 3:1 downsampling. - audioop.ratecv produces muffled audio because it lacks an anti-aliasing filter. - """ - # Downsample from 24kHz to 8kHz using high-quality resampler - audio_data = np.frombuffer(pcm_bytes, dtype=np.int16) - resampled = soxr.resample(audio_data, 24000, 8000, quality="VHQ") - # Both audioop.ratecv (upstream) and soxr can produce ±1 sample due to filter rounding. - # Use round() so that e.g. 2399 input samples → round(2399/3) = 800, not 799. - expected_samples = round(len(audio_data) * 8000 / 24000) - if len(resampled) < expected_samples: - resampled = np.pad(resampled, (0, expected_samples - len(resampled))) - elif len(resampled) > expected_samples: - resampled = resampled[:expected_samples] - pcm_8k = resampled.astype(np.int16).tobytes() - # Encode to mu-law - return audioop.lin2ulaw(pcm_8k, 2) - - -def sync_buffer_to_position(buffer: bytearray, target_position: int) -> None: - """Pad *buffer* with silence bytes so it reaches *target_position*. - - Mirrors pipecat's ``AudioBufferProcessor._sync_buffer_to_position``. - Call this **before** extending the *other* track so both tracks stay - positionally aligned. - """ - current_len = len(buffer) - if current_len < target_position: - buffer.extend(b"\x00" * (target_position - current_len)) - - -def pcm16_mix(track_a: bytes, track_b: bytes) -> bytes: - """Mix two 16-bit PCM tracks by sample-wise addition with clipping. - - Both tracks must be the same sample rate. If lengths differ, - the shorter track is zero-padded. - """ - len_a, len_b = len(track_a), len(track_b) - max_len = max(len_a, len_b) - - # Zero-pad shorter track - if len_a < max_len: - track_a = track_a + b"\x00" * (max_len - len_a) - if len_b < max_len: - track_b = track_b + b"\x00" * (max_len - len_b) - - # Mix with clipping - n_samples = max_len // 2 - fmt = f"<{n_samples}h" - samples_a = struct.unpack(fmt, track_a) - samples_b = struct.unpack(fmt, track_b) - mixed = struct.pack(fmt, *(max(-32768, min(32767, a + b)) for a, b in zip(samples_a, samples_b))) - return mixed - - -# ── Twilio WebSocket Protocol ──────────────────────────────────────── - - -def parse_twilio_media_message(message: str) -> bytes | None: - """Parse a Twilio media WebSocket message and extract raw audio bytes. - - Returns None if the message is not a media message. - """ - try: - data = json.loads(message) - if data.get("event") == "media": - payload = data["media"]["payload"] - return base64.b64decode(payload) - except (json.JSONDecodeError, KeyError): - pass - return None - - -def create_twilio_media_message(stream_sid: str, audio_bytes: bytes) -> str: - """Create a Twilio media WebSocket message with the given audio bytes.""" - payload = base64.b64encode(audio_bytes).decode("ascii") - return json.dumps( - { - "event": "media", - "streamSid": stream_sid, - "media": { - "payload": payload, - }, - } - ) - - -# ── Framework Logs Writer ──────────────────────────────────────────── - - -class FrameworkLogWriter: - """Write framework_logs.jsonl (replacement for pipecat_logs.jsonl). - - Capture turn boundaries, TTS text, and LLM responses with accurate - wall-clock timestamps. - """ - - def __init__(self, output_dir: Path): - self.log_file = output_dir / "framework_logs.jsonl" - output_dir.mkdir(parents=True, exist_ok=True) - - def write(self, event_type: str, data: dict, timestamp_ms: int | None = None) -> None: - """Write a single log entry. - - Args: - event_type: One of 'turn_start', 'turn_end', 'tts_text', 'llm_response' - data: Event data dict. Must contain a 'frame' key for tts_text/llm_response. - timestamp_ms: Wall-clock timestamp in milliseconds. Defaults to now. - """ - if timestamp_ms is None: - timestamp_ms = int(time.time() * 1000) - - entry = { - "timestamp": timestamp_ms, - "type": event_type, - "data": data, - } - try: - with open(self.log_file, "a", encoding="utf-8") as f: - f.write(json.dumps(entry, ensure_ascii=False) + "\n") - except Exception as e: - logger.error(f"Error writing framework log: {e}") - - def turn_start(self, timestamp_ms: int | None = None) -> None: - """Log a turn start event.""" - self.write("turn_start", {"frame": "turn_start"}, timestamp_ms) - - def turn_end(self, was_interrupted: bool = False, timestamp_ms: int | None = None) -> None: - """Log a turn end event.""" - self.write("turn_end", {"frame": "turn_end", "was_interrupted": was_interrupted}, timestamp_ms) - - def tts_text(self, text: str, timestamp_ms: int | None = None) -> None: - """Log TTS text (what was actually spoken).""" - self.write("tts_text", {"frame": text}, timestamp_ms) - - def llm_response(self, text: str, timestamp_ms: int | None = None) -> None: - """Log LLM response text (full intended response).""" - self.write("llm_response", {"frame": text}, timestamp_ms) - - def s2s_transcript(self, text: str, timestamp_ms: int | None = None) -> None: - """Log S2S transcript (what was actually spoken).""" - self.write("s2s_transcript", {"frame": text}, timestamp_ms) - - -# ── Metrics Log Writer ─────────────────────────────────────────────── - - -class MetricsLogWriter: - """Writes pipecat_metrics.jsonl for non-Pipecat frameworks. - - Pipecat writes its own metrics natively via MetricsFileObserver. This - writer covers OpenAI Realtime, Gemini Live, and any other framework that - manages its own session loop. - """ - - def __init__(self, output_dir: Path): - self.log_file = output_dir / "pipecat_metrics.jsonl" - output_dir.mkdir(parents=True, exist_ok=True) - - def write_latency(self, stage: str, value_seconds: float, model: str = "") -> None: - """Write a LatencyMetric entry. - - Args: - stage: Semantic label for the stage being measured. Use ``"stt"`` - for STT processing time, ``"tts"`` for TTS time-to-first-byte, - ``"model_response"`` for s2s/realtime time from user speech end - to first model audio chunk. - value_seconds: Latency in seconds. - model: Model identifier (optional). - """ - entry = { - "timestamp": int(time.time() * 1000), - "type": "LatencyMetric", - "stage": stage, - "model": model, - "value": value_seconds, - } - self._append(entry) - - def write_token_usage( - self, - processor: str, - model: str, - prompt_tokens: int, - completion_tokens: int, - ) -> None: - """Write an LLMTokenUsageMetricsData entry.""" - entry = { - "timestamp": int(time.time() * 1000), - "type": "LLMTokenUsageMetricsData", - "processor": processor, - "model": model, - "value": { - "prompt_tokens": prompt_tokens, - "completion_tokens": completion_tokens, - "total_tokens": prompt_tokens + completion_tokens, - }, - } - self._append(entry) - - def _append(self, entry: dict) -> None: - try: - with open(self.log_file, "a", encoding="utf-8") as f: - f.write(json.dumps(entry, ensure_ascii=False) + "\n") - except Exception as e: - logger.error(f"Error writing metrics log: {e}") diff --git a/src/eva/assistant/base_server.py b/src/eva/assistant/base_server.py index 43371fce..4fa4886d 100644 --- a/src/eva/assistant/base_server.py +++ b/src/eva/assistant/base_server.py @@ -16,11 +16,11 @@ from fastapi import FastAPI from eva.assistant.agentic.audit_log import AuditLog -from eva.assistant.audio_bridge import FrameworkLogWriter, MetricsLogWriter +from eva.assistant.pipeline.observers import FrameworkLogWriter, MetricsLogWriter from eva.assistant.tools.tool_executor import ToolExecutor from eva.models.agents import AgentConfig from eva.models.config import ModelConfig -from eva.utils.audio_utils import save_pcm_as_wav +from eva.utils.audio_utils import pcm16_mix, save_pcm_as_wav from eva.utils.culture import get_initial_message from eva.utils.logging import get_logger from eva.utils.prompt_manager import PromptManager @@ -166,8 +166,6 @@ async def stop(self) -> asyncio.Task | None: f"assistant={len(self.assistant_audio_buffer)} " f"diff={diff_ms:.0f}ms — mixed recording may be temporally skewed" ) - from eva.assistant.audio_bridge import pcm16_mix # lazy: avoids circular import at module load - self._audio_buffer = bytearray( pcm16_mix(bytes(self.user_audio_buffer), bytes(self.assistant_audio_buffer)) ) @@ -291,8 +289,6 @@ def _save_audio(self) -> None: f"assistant={len(self.assistant_audio_buffer)} " f"diff={diff_ms:.0f}ms — mixed recording may be temporally skewed" ) - from eva.assistant.audio_bridge import pcm16_mix - self._audio_buffer = bytearray(pcm16_mix(bytes(self.user_audio_buffer), bytes(self.assistant_audio_buffer))) elif not self._audio_buffer and self.user_audio_buffer: self._audio_buffer = bytearray(self.user_audio_buffer) diff --git a/src/eva/assistant/elevenlabs_server.py b/src/eva/assistant/elevenlabs_server.py index 0a543440..8d37e276 100644 --- a/src/eva/assistant/elevenlabs_server.py +++ b/src/eva/assistant/elevenlabs_server.py @@ -33,17 +33,16 @@ ) from fastapi import FastAPI, WebSocket, WebSocketDisconnect -from eva.assistant.audio_bridge import ( - FrameworkLogWriter, - MetricsLogWriter, - create_twilio_media_message, - mulaw_8k_to_pcm16_16k, - parse_twilio_media_message, -) from eva.assistant.base_server import AbstractAssistantServer from eva.assistant.elevenlabs_audio_interface import TwilioAudioBridge +from eva.assistant.pipeline.observers import FrameworkLogWriter, MetricsLogWriter from eva.models.agents import AgentConfig from eva.models.config import ModelConfig +from eva.utils.audio_utils import ( + create_twilio_media_message, + mulaw_8k_to_pcm16_16k, + parse_twilio_media_message, +) from eva.utils.logging import get_logger logger = get_logger(__name__) diff --git a/src/eva/assistant/gemini_live_server.py b/src/eva/assistant/gemini_live_server.py index de69e32c..377386a9 100644 --- a/src/eva/assistant/gemini_live_server.py +++ b/src/eva/assistant/gemini_live_server.py @@ -26,9 +26,11 @@ from google import genai from google.genai import types -from eva.assistant.audio_bridge import ( - FrameworkLogWriter, - MetricsLogWriter, +from eva.assistant.base_server import AbstractAssistantServer +from eva.assistant.pipeline.observers import FrameworkLogWriter, MetricsLogWriter +from eva.models.agents import AgentConfig +from eva.models.config import ModelConfig +from eva.utils.audio_utils import ( create_twilio_media_message, mulaw_8k_to_pcm16_16k, mulaw_8k_to_pcm16_24k, @@ -36,9 +38,6 @@ pcm16_24k_to_mulaw_8k, sync_buffer_to_position, ) -from eva.assistant.base_server import AbstractAssistantServer -from eva.models.agents import AgentConfig -from eva.models.config import ModelConfig from eva.utils.logging import get_logger logger = get_logger(__name__) diff --git a/src/eva/assistant/openai_realtime_server.py b/src/eva/assistant/openai_realtime_server.py index 27661537..0c29f16f 100644 --- a/src/eva/assistant/openai_realtime_server.py +++ b/src/eva/assistant/openai_realtime_server.py @@ -17,16 +17,15 @@ from fastapi import FastAPI, WebSocket, WebSocketDisconnect from openai import AsyncOpenAI -from eva.assistant.audio_bridge import ( - FrameworkLogWriter, - MetricsLogWriter, +from eva.assistant.base_server import AbstractAssistantServer +from eva.assistant.pipeline.observers import FrameworkLogWriter, MetricsLogWriter +from eva.utils.audio_utils import ( create_twilio_media_message, mulaw_8k_to_pcm16_24k, parse_twilio_media_message, pcm16_24k_to_mulaw_8k, sync_buffer_to_position, ) -from eva.assistant.base_server import AbstractAssistantServer from eva.utils.logging import get_logger logger = get_logger(__name__) diff --git a/src/eva/assistant/pipeline/observers.py b/src/eva/assistant/pipeline/observers.py index 90c035fb..a84f20f5 100644 --- a/src/eva/assistant/pipeline/observers.py +++ b/src/eva/assistant/pipeline/observers.py @@ -348,3 +348,105 @@ def close(self) -> None: def __del__(self): """Fallback cleanup — prefer calling close() explicitly.""" self.close() + + +# ── Non-Pipecat log writers ────────────────────────────────────────── +# +# The observers above are Pipecat-native. The writers below are their +# equivalents for frameworks that manage their own session loop (OpenAI +# Realtime, Gemini Live, ElevenLabs) and therefore can't use the observers. + + +class FrameworkLogWriter: + """Write framework_logs.jsonl (replacement for pipecat_logs.jsonl). + + Capture turn boundaries, TTS text, and LLM responses with accurate + wall-clock timestamps. + """ + + def __init__(self, output_dir: Path): + self.log_file = output_dir / "framework_logs.jsonl" + output_dir.mkdir(parents=True, exist_ok=True) + + def write(self, event_type: str, data: dict, timestamp_ms: int | None = None) -> None: + """Write a single log entry. + + Args: + event_type: One of 'turn_start', 'turn_end', 'tts_text', 'llm_response' + data: Event data dict. Must contain a 'frame' key for tts_text/llm_response. + timestamp_ms: Wall-clock timestamp in milliseconds. Defaults to now. + """ + if timestamp_ms is None: + timestamp_ms = int(time.time() * 1000) + + entry = { + "timestamp": timestamp_ms, + "type": event_type, + "data": data, + } + try: + with open(self.log_file, "a", encoding="utf-8") as f: + f.write(json.dumps(entry, ensure_ascii=False) + "\n") + except Exception as e: + logger.error(f"Error writing framework log: {e}") + + +class MetricsLogWriter: + """Writes pipecat_metrics.jsonl for non-Pipecat frameworks. + + Pipecat writes its own metrics natively via MetricsFileObserver. This + writer covers OpenAI Realtime, Gemini Live, and any other framework that + manages its own session loop. + """ + + def __init__(self, output_dir: Path): + self.log_file = output_dir / "pipecat_metrics.jsonl" + output_dir.mkdir(parents=True, exist_ok=True) + + def write_latency(self, stage: str, value_seconds: float, model: str = "") -> None: + """Write a LatencyMetric entry. + + Args: + stage: Semantic label for the stage being measured. Use ``"stt"`` + for STT processing time, ``"tts"`` for TTS time-to-first-byte, + ``"model_response"`` for s2s/realtime time from user speech end + to first model audio chunk. + value_seconds: Latency in seconds. + model: Model identifier (optional). + """ + entry = { + "timestamp": int(time.time() * 1000), + "type": "LatencyMetric", + "stage": stage, + "model": model, + "value": value_seconds, + } + self._append(entry) + + def write_token_usage( + self, + processor: str, + model: str, + prompt_tokens: int, + completion_tokens: int, + ) -> None: + """Write an LLMTokenUsageMetricsData entry.""" + entry = { + "timestamp": int(time.time() * 1000), + "type": "LLMTokenUsageMetricsData", + "processor": processor, + "model": model, + "value": { + "prompt_tokens": prompt_tokens, + "completion_tokens": completion_tokens, + "total_tokens": prompt_tokens + completion_tokens, + }, + } + self._append(entry) + + def _append(self, entry: dict) -> None: + try: + with open(self.log_file, "a", encoding="utf-8") as f: + f.write(json.dumps(entry, ensure_ascii=False) + "\n") + except Exception as e: + logger.error(f"Error writing metrics log: {e}") diff --git a/src/eva/utils/audio_utils.py b/src/eva/utils/audio_utils.py index 77c82eef..ad38c37d 100644 --- a/src/eva/utils/audio_utils.py +++ b/src/eva/utils/audio_utils.py @@ -1,13 +1,140 @@ -"""Shared audio I/O helpers.""" +"""Shared audio I/O and format-conversion helpers.""" +import audioop +import base64 +import json +import struct import wave from pathlib import Path +import numpy as np +import soxr + from eva.utils.logging import get_logger logger = get_logger(__name__) +# ── Audio format conversion ────────────────────────────────────────── + + +def mulaw_8k_to_pcm16_16k(mulaw_bytes: bytes) -> bytes: + """Convert 8kHz mu-law audio to 16kHz 16-bit PCM.""" + # Decode mu-law to 16-bit PCM at 8kHz + pcm_8k = audioop.ulaw2lin(mulaw_bytes, 2) + # Upsample from 8kHz to 16kHz + pcm_16k, _ = audioop.ratecv(pcm_8k, 2, 1, 8000, 16000, None) + return pcm_16k + + +def mulaw_8k_to_pcm16_24k(mulaw_bytes: bytes) -> bytes: + """Convert 8kHz mu-law audio to 24kHz 16-bit PCM.""" + # Decode mu-law to 16-bit PCM at 8kHz + pcm_8k = audioop.ulaw2lin(mulaw_bytes, 2) + # Upsample from 8kHz to 24kHz + pcm_24k, _ = audioop.ratecv(pcm_8k, 2, 1, 8000, 24000, None) + # audioop.ratecv can produce ±2 samples; clamp to exact 3× input length + # so that the inverse conversion recovers the original sample count. + expected_bytes = len(pcm_8k) * 3 + if len(pcm_24k) < expected_bytes: + pcm_24k = pcm_24k + b"\x00" * (expected_bytes - len(pcm_24k)) + elif len(pcm_24k) > expected_bytes: + pcm_24k = pcm_24k[:expected_bytes] + return pcm_24k + + +def pcm16_24k_to_mulaw_8k(pcm_bytes: bytes) -> bytes: + """Convert 24kHz 16-bit PCM to 8kHz mu-law. + + Uses soxr VHQ resampling (same as Pipecat) for proper anti-aliasing during the 3:1 downsampling. + audioop.ratecv produces muffled audio because it lacks an anti-aliasing filter. + """ + # Downsample from 24kHz to 8kHz using high-quality resampler + audio_data = np.frombuffer(pcm_bytes, dtype=np.int16) + resampled = soxr.resample(audio_data, 24000, 8000, quality="VHQ") + # Both audioop.ratecv (upstream) and soxr can produce ±1 sample due to filter rounding. + # Use round() so that e.g. 2399 input samples → round(2399/3) = 800, not 799. + expected_samples = round(len(audio_data) * 8000 / 24000) + if len(resampled) < expected_samples: + resampled = np.pad(resampled, (0, expected_samples - len(resampled))) + elif len(resampled) > expected_samples: + resampled = resampled[:expected_samples] + pcm_8k = resampled.astype(np.int16).tobytes() + # Encode to mu-law + return audioop.lin2ulaw(pcm_8k, 2) + + +def sync_buffer_to_position(buffer: bytearray, target_position: int) -> None: + """Pad *buffer* with silence bytes so it reaches *target_position*. + + Mirrors pipecat's ``AudioBufferProcessor._sync_buffer_to_position``. + Call this **before** extending the *other* track so both tracks stay + positionally aligned. + """ + current_len = len(buffer) + if current_len < target_position: + buffer.extend(b"\x00" * (target_position - current_len)) + + +def pcm16_mix(track_a: bytes, track_b: bytes) -> bytes: + """Mix two 16-bit PCM tracks by sample-wise addition with clipping. + + Both tracks must be the same sample rate. If lengths differ, + the shorter track is zero-padded. + """ + len_a, len_b = len(track_a), len(track_b) + max_len = max(len_a, len_b) + + # Zero-pad shorter track + if len_a < max_len: + track_a = track_a + b"\x00" * (max_len - len_a) + if len_b < max_len: + track_b = track_b + b"\x00" * (max_len - len_b) + + # Mix with clipping + n_samples = max_len // 2 + fmt = f"<{n_samples}h" + samples_a = struct.unpack(fmt, track_a) + samples_b = struct.unpack(fmt, track_b) + mixed = struct.pack(fmt, *(max(-32768, min(32767, a + b)) for a, b in zip(samples_a, samples_b))) + return mixed + + +# ── Twilio WebSocket Protocol ──────────────────────────────────────── + + +def parse_twilio_media_message(message: str) -> bytes | None: + """Parse a Twilio media WebSocket message and extract raw audio bytes. + + Returns None if the message is not a media message. + """ + try: + data = json.loads(message) + if data.get("event") == "media": + payload = data["media"]["payload"] + return base64.b64decode(payload) + except (json.JSONDecodeError, KeyError): + pass + return None + + +def create_twilio_media_message(stream_sid: str, audio_bytes: bytes) -> str: + """Create a Twilio media WebSocket message with the given audio bytes.""" + payload = base64.b64encode(audio_bytes).decode("ascii") + return json.dumps( + { + "event": "media", + "streamSid": stream_sid, + "media": { + "payload": payload, + }, + } + ) + + +# ── WAV I/O ────────────────────────────────────────────────────────── + + def save_pcm_as_wav( audio_data: bytes, file_path: Path, diff --git a/tests/unit/assistant/test_audio_bridge.py b/tests/unit/utils/test_audio_utils.py similarity index 95% rename from tests/unit/assistant/test_audio_bridge.py rename to tests/unit/utils/test_audio_utils.py index 211c68a6..748a154f 100644 --- a/tests/unit/assistant/test_audio_bridge.py +++ b/tests/unit/utils/test_audio_utils.py @@ -1,4 +1,4 @@ -"""Tests for shared audio bridge utilities. +"""Tests for shared audio format-conversion utilities in eva.utils.audio_utils. Covers: PCM↔mulaw round-trip fidelity, PCM16 mixing with clipping, and Twilio WebSocket protocol message round-trips. @@ -11,7 +11,7 @@ import pytest -from eva.assistant.audio_bridge import ( +from eva.utils.audio_utils import ( create_twilio_media_message, mulaw_8k_to_pcm16_24k, parse_twilio_media_message,