Use push_error in services

This commit is contained in:
Mark Backman
2025-11-24 18:10:08 -05:00
parent a475d09b90
commit c6fcf58be4
50 changed files with 156 additions and 230 deletions

View File

@@ -285,12 +285,11 @@ class AsyncAITTSService(InterruptibleTTSService):
) )
await self.push_frame(frame) await self.push_frame(frame)
elif msg.get("error_code"): elif msg.get("error_code"):
logger.error(f"{self} error: {msg}")
await self.push_frame(TTSStoppedFrame()) await self.push_frame(TTSStoppedFrame())
await self.stop_all_metrics() await self.stop_all_metrics()
await self.push_error(error_msg=f"{self} error: {msg['message']}") await self.push_error(error_msg=f"{self} error: {msg['message']}")
else: else:
logger.error(f"{self} error, unknown message type: {msg}") await self.push_error(error_msg=f"{self} error, unknown message type: {msg}")
async def _keepalive_task_handler(self): async def _keepalive_task_handler(self):
"""Send periodic keepalive messages to maintain WebSocket connection.""" """Send periodic keepalive messages to maintain WebSocket connection."""
@@ -333,16 +332,14 @@ class AsyncAITTSService(InterruptibleTTSService):
await self._get_websocket().send(msg) await self._get_websocket().send(msg)
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class AsyncAIHttpTTSService(TTSService): class AsyncAIHttpTTSService(TTSService):

View File

@@ -452,7 +452,7 @@ class AWSNovaSonicLLMService(LLMService):
self._ready_to_send_context = True self._ready_to_send_context = True
await self._finish_connecting_if_context_available() await self._finish_connecting_if_context_available()
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(exception=e)
await self._disconnect() await self._disconnect()
async def _process_completed_function_calls(self, send_new_results: bool): async def _process_completed_function_calls(self, send_new_results: bool):
@@ -576,7 +576,7 @@ class AWSNovaSonicLLMService(LLMService):
logger.info("Finished disconnecting") logger.info("Finished disconnecting")
except Exception as e: except Exception as e:
logger.error(f"{self} error disconnecting: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
def _create_client(self) -> BedrockRuntimeClient: def _create_client(self) -> BedrockRuntimeClient:
config = Config( config = Config(
@@ -884,7 +884,7 @@ class AWSNovaSonicLLMService(LLMService):
# Errors are kind of expected while disconnecting, so just # Errors are kind of expected while disconnecting, so just
# ignore them and do nothing # ignore them and do nothing
return return
logger.error(f"{self} error processing responses: {e}") await self.push_error(exception=e)
if self._wants_connection: if self._wants_connection:
await self.reset_conversation() await self.reset_conversation()

View File

@@ -181,8 +181,7 @@ class AWSTranscribeSTTService(STTService):
try: try:
await self._connect() await self._connect()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
return return
# Format the audio data according to AWS event stream format # Format the audio data according to AWS event stream format
@@ -199,13 +198,11 @@ class AWSTranscribeSTTService(STTService):
await self._disconnect() await self._disconnect()
# Don't yield error here - we'll retry on next frame # Don't yield error here - we'll retry on next frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
await self._disconnect() await self._disconnect()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
await self._disconnect() await self._disconnect()
async def _connect(self): async def _connect(self):
@@ -526,13 +523,14 @@ class AWSTranscribeSTTService(STTService):
) )
elif headers.get(":message-type") == "exception": elif headers.get(":message-type") == "exception":
error_msg = payload.get("Message", "Unknown error") error_msg = payload.get("Message", "Unknown error")
logger.error(f"{self} Exception from AWS: {error_msg}") await self.push_error(error_msg=f"AWS Transcribe error: {error_msg}")
await self.push_frame(ErrorFrame(f"AWS Transcribe error: {error_msg}"))
else: else:
logger.debug(f"{self} Other message type received: {headers}") logger.debug(f"{self} Other message type received: {headers}")
logger.debug(f"{self} Payload: {payload}") logger.debug(f"{self} Payload: {payload}")
except websockets.exceptions.ConnectionClosed as e: except websockets.exceptions.ConnectionClosed as e:
logger.error(f"{self} WebSocket connection closed in receive loop: {e}") await self.push_error(
error_msg=f"WebSocket connection closed in receive loop", exception=e
)
break break
except Exception as e: except Exception as e:
await self.push_error(exception=e) await self.push_error(exception=e)

View File

@@ -311,10 +311,8 @@ class AWSPollyTTSService(TTSService):
yield frame yield frame
yield TTSStoppedFrame() yield TTSStoppedFrame()
except (BotoCoreError, ClientError) as error: except (BotoCoreError, ClientError) as e:
logger.error(f"{self} error generating TTS: {error}") await self.push_error(exception=e)
error_message = f"AWS Polly TTS error: {str(error)}"
yield ErrorFrame(error=error_message)
finally: finally:
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -91,8 +91,7 @@ class AzureImageGenServiceREST(ImageGenService):
while status != "succeeded": while status != "succeeded":
attempts_left -= 1 attempts_left -= 1
if attempts_left == 0: if attempts_left == 0:
logger.error(f"{self} error: image generation timed out") await self.push_error(error_msg="Image generation timed out")
yield ErrorFrame("Image generation timed out")
return return
await asyncio.sleep(1) await asyncio.sleep(1)
@@ -104,8 +103,7 @@ class AzureImageGenServiceREST(ImageGenService):
image_url = json_response["result"]["data"][0]["url"] if json_response else None image_url = json_response["result"]["data"][0]["url"] if json_response else None
if not image_url: if not image_url:
logger.error(f"{self} error: image generation failed") await self.push_error(error_msg="Image generation failed")
yield ErrorFrame("Image generation failed")
return return
# Load the image from the url # Load the image from the url

View File

@@ -61,5 +61,5 @@ class AzureRealtimeLLMService(OpenAIRealtimeLLMService):
) )
self._receive_task = self.create_task(self._receive_task_handler()) self._receive_task = self.create_task(self._receive_task_handler())
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(exception=e)
self._websocket = None self._websocket = None

View File

@@ -121,8 +121,7 @@ class AzureSTTService(STTService):
self._audio_stream.write(audio) self._audio_stream.write(audio)
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
async def start(self, frame: StartFrame): async def start(self, frame: StartFrame):
"""Start the speech recognition service. """Start the speech recognition service.

View File

@@ -327,8 +327,7 @@ class AzureTTSService(AzureBaseTTSService):
try: try:
if self._speech_synthesizer is None: if self._speech_synthesizer is None:
error_msg = "Speech synthesizer not initialized." error_msg = "Speech synthesizer not initialized."
logger.error(error_msg) await self.push_error(error_msg=error_msg)
yield ErrorFrame(error=error_msg)
return return
try: try:
@@ -355,15 +354,13 @@ class AzureTTSService(AzureBaseTTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
# Could add reconnection logic here if needed # Could add reconnection logic here if needed
return return
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class AzureHttpTTSService(AzureBaseTTSService): class AzureHttpTTSService(AzureBaseTTSService):
@@ -440,5 +437,4 @@ class AzureHttpTTSService(AzureBaseTTSService):
cancellation_details = result.cancellation_details cancellation_details = result.cancellation_details
logger.warning(f"Speech synthesis canceled: {cancellation_details.reason}") logger.warning(f"Speech synthesis canceled: {cancellation_details.reason}")
if cancellation_details.reason == CancellationReason.Error: if cancellation_details.reason == CancellationReason.Error:
logger.error(f"{self} error: {cancellation_details.error_details}") await self.push_error(error_msg=cancellation_details.error_details)
yield ErrorFrame(error=f"{self} error: {cancellation_details.error_details}")

View File

@@ -467,7 +467,7 @@ class CartesiaTTSService(AudioContextWordTTSService):
await self.push_error(error_msg=f"{self} error: {msg}") await self.push_error(error_msg=f"{self} error: {msg}")
self._context_id = None self._context_id = None
else: else:
logger.error(f"{self} error, unknown message type: {msg}") await self.push_error(error_msg=f"{self} error, unknown message type: {msg}")
async def _receive_messages(self): async def _receive_messages(self):
while True: while True:
@@ -505,16 +505,14 @@ class CartesiaTTSService(AudioContextWordTTSService):
await self._get_websocket().send(msg) await self._get_websocket().send(msg)
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class CartesiaHttpTTSService(TTSService): class CartesiaHttpTTSService(TTSService):

View File

@@ -378,16 +378,14 @@ class DeepgramFluxSTTService(WebsocketSTTService):
are issues sending the audio data. are issues sending the audio data.
""" """
if not self._websocket: if not self._websocket:
logger.error("Not connected to Deepgram Flux.") await self.push_error(error_msg=f"Websocket not connected")
yield ErrorFrame("Not connected to Deepgram Flux.")
return return
try: try:
self._last_stt_time = time.monotonic() self._last_stt_time = time.monotonic()
await self.send_with_retry(audio, self._report_error) await self.send_with_retry(audio, self._report_error)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
return return
yield None yield None

View File

@@ -233,7 +233,7 @@ class DeepgramSTTService(STTService):
) )
if not await self._connection.start(options=self._settings, addons=self._addons): if not await self._connection.start(options=self._settings, addons=self._addons):
logger.error(f"{self}: unable to connect to Deepgram") await self.push_error(error_msg=f"Unable to connect to Deepgram")
async def _disconnect(self): async def _disconnect(self):
if await self._connection.is_connected(): if await self._connection.is_connected():

View File

@@ -116,8 +116,7 @@ class DeepgramTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class DeepgramHttpTTSService(TTSService): class DeepgramHttpTTSService(TTSService):
@@ -227,5 +226,4 @@ class DeepgramHttpTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(f"Error getting audio: {str(e)}")

View File

@@ -351,8 +351,7 @@ class ElevenLabsSTTService(SegmentedSTTService):
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
def audio_format_from_sample_rate(sample_rate: int) -> str: def audio_format_from_sample_rate(sample_rate: int) -> str:
@@ -586,8 +585,7 @@ class ElevenLabsRealtimeSTTService(WebsocketSTTService):
} }
await self._websocket.send(json.dumps(message)) await self._websocket.send(json.dumps(message))
except Exception as e: except Exception as e:
logger.error(f"Error sending audio: {e}") await self.push_error(exception=e)
yield ErrorFrame(f"ElevenLabs Realtime STT error: {str(e)}")
yield None yield None
@@ -656,7 +654,7 @@ class ElevenLabsRealtimeSTTService(WebsocketSTTService):
logger.debug("Disconnecting from ElevenLabs Realtime STT") logger.debug("Disconnecting from ElevenLabs Realtime STT")
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} error closing websocket: {e}") await self.push_error(error_msg=f"Error closing websocket: {e}", exception=e)
finally: finally:
self._websocket = None self._websocket = None
await self._call_event_handler("on_disconnected") await self._call_event_handler("on_disconnected")

View File

@@ -535,9 +535,8 @@ class ElevenLabsTTSService(AudioContextWordTTSService):
await self._call_event_handler("on_connected") await self._call_event_handler("on_connected")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}")
self._websocket = None self._websocket = None
await self.push_error(ErrorFrame(error=f"{self} error: {e}")) await self.push_error(exception=e)
await self._call_event_handler("on_connection_error", f"{e}") await self._call_event_handler("on_connection_error", f"{e}")
async def _disconnect_websocket(self): async def _disconnect_websocket(self):
@@ -582,8 +581,7 @@ class ElevenLabsTTSService(AudioContextWordTTSService):
json.dumps({"context_id": self._context_id, "close_context": True}) json.dumps({"context_id": self._context_id, "close_context": True})
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
self._context_id = None self._context_id = None
self._started = False self._started = False
self._partial_word = "" self._partial_word = ""
@@ -738,15 +736,13 @@ class ElevenLabsTTSService(AudioContextWordTTSService):
else: else:
await self._send_text(text) await self._send_text(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield TTSStoppedFrame() yield TTSStoppedFrame()
yield ErrorFrame(error=f"{self} error: {e}")
self._started = False self._started = False
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class ElevenLabsHttpTTSService(WordTTSService): class ElevenLabsHttpTTSService(WordTTSService):
@@ -1041,8 +1037,7 @@ class ElevenLabsHttpTTSService(WordTTSService):
) as response: ) as response:
if response.status != 200: if response.status != 200:
error_text = await response.text() error_text = await response.text()
logger.error(f"{self} error: {error_text}") await self.push_error(error_msg=f"ElevenLabs API error: {error_text}")
yield ErrorFrame(error=f"ElevenLabs API error: {error_text}")
return return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
@@ -1089,8 +1084,7 @@ class ElevenLabsHttpTTSService(WordTTSService):
logger.warning(f"Failed to parse JSON from stream: {e}") logger.warning(f"Failed to parse JSON from stream: {e}")
continue continue
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
continue continue
# After processing all chunks, emit any remaining partial word # After processing all chunks, emit any remaining partial word
@@ -1114,8 +1108,7 @@ class ElevenLabsHttpTTSService(WordTTSService):
self._previous_text = text self._previous_text = text
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
# Let the parent class handle TTSStoppedFrame # Let the parent class handle TTSStoppedFrame

View File

@@ -110,8 +110,7 @@ class FalImageGenService(ImageGenService):
image_url = response["images"][0]["url"] if response else None image_url = response["images"][0]["url"] if response else None
if not image_url: if not image_url:
logger.error(f"{self} error: image generation failed") await self.push_error(error_msg="Image generation failed")
yield ErrorFrame("Image generation failed")
return return
logger.debug(f"Image generated at: {image_url}") logger.debug(f"Image generated at: {image_url}")

View File

@@ -290,5 +290,4 @@ class FalSTTService(SegmentedSTTService):
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -320,8 +320,7 @@ class FishAudioTTSService(InterruptibleTTSService):
flush_message = {"event": "flush"} flush_message = {"event": "flush"}
await self._get_websocket().send(ormsgpack.packb(flush_message)) await self._get_websocket().send(ormsgpack.packb(flush_message))
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
@@ -329,5 +328,4 @@ class FishAudioTTSService(InterruptibleTTSService):
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -558,8 +558,7 @@ class GladiaSTTService(STTService):
except websockets.exceptions.ConnectionClosed: except websockets.exceptions.ConnectionClosed:
logger.debug("Connection closed during keepalive") logger.debug("Connection closed during keepalive")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def _receive_task_handler(self): async def _receive_task_handler(self):
try: try:
@@ -630,7 +629,10 @@ class GladiaSTTService(STTService):
return False return False
self._reconnection_attempts += 1 self._reconnection_attempts += 1
if self._reconnection_attempts > self._max_reconnection_attempts: if self._reconnection_attempts > self._max_reconnection_attempts:
logger.error(f"Max reconnection attempts ({self._max_reconnection_attempts}) reached") await self.push_error(
error_msg=f"Max reconnection attempts ({self._max_reconnection_attempts}) reached",
fatal=True,
)
self._should_reconnect = False self._should_reconnect = False
return False return False
delay = self._reconnection_delay * (2 ** (self._reconnection_attempts - 1)) delay = self._reconnection_delay * (2 ** (self._reconnection_attempts - 1))

View File

@@ -1251,11 +1251,11 @@ class GeminiLiveLLMService(LLMService):
) )
if self._consecutive_failures >= MAX_CONSECUTIVE_FAILURES: if self._consecutive_failures >= MAX_CONSECUTIVE_FAILURES:
logger.error( error_msg = (
f"Max consecutive failures ({MAX_CONSECUTIVE_FAILURES}) reached, " f"Max consecutive failures ({MAX_CONSECUTIVE_FAILURES}) reached, "
"treating as fatal error" "treating as fatal error"
) )
await self.push_error(ErrorFrame(error=f"{self} Error in receive loop: {error}")) await self.push_error(error_msg=error_msg, exception=error, fatal=True)
return False return False
else: else:
logger.info( logger.info(
@@ -1283,7 +1283,7 @@ class GeminiLiveLLMService(LLMService):
self._completed_tool_calls = set() self._completed_tool_calls = set()
self._disconnecting = False self._disconnecting = False
except Exception as e: except Exception as e:
logger.error(f"{self} error disconnecting: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
async def _send_user_audio(self, frame): async def _send_user_audio(self, frame):
"""Send user audio frame to Gemini Live API.""" """Send user audio frame to Gemini Live API."""

View File

@@ -110,8 +110,7 @@ class GoogleImageGenService(ImageGenService):
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
if not response or not response.generated_images: if not response or not response.generated_images:
logger.error(f"{self} error: image generation failed") await self.push_error(error_msg="Image generation failed")
yield ErrorFrame("Image generation failed")
return return
for img_response in response.generated_images: for img_response in response.generated_images:
@@ -128,5 +127,4 @@ class GoogleImageGenService(ImageGenService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} error generating image: {e}") await self.push_error(exception=e)
yield ErrorFrame(f"Image generation error: {str(e)}")

View File

@@ -983,7 +983,7 @@ class GoogleLLMService(LLMService):
except DeadlineExceeded: except DeadlineExceeded:
await self._call_event_handler("on_completion_timeout") await self._call_event_handler("on_completion_timeout")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
finally: finally:
if grounding_metadata and isinstance(grounding_metadata, dict): if grounding_metadata and isinstance(grounding_metadata, dict):
llm_search_frame = LLMSearchResponseFrame( llm_search_frame = LLMSearchResponseFrame(

View File

@@ -805,8 +805,7 @@ class GoogleSTTService(STTService):
break break
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
await asyncio.sleep(1) # Brief delay before reconnecting await asyncio.sleep(1) # Brief delay before reconnecting
self._stream_start_time = int(time.time() * 1000) self._stream_start_time = int(time.time() * 1000)
@@ -900,8 +899,7 @@ class GoogleSTTService(STTService):
) )
raise raise
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
# Re-raise the exception to let it propagate (e.g. in the case of a # Re-raise the exception to let it propagate (e.g. in the case of a
# timeout, propagate to _stream_audio to reconnect) # timeout, propagate to _stream_audio to reconnect)
raise raise

View File

@@ -737,9 +737,7 @@ class GoogleHttpTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
error_message = f"TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message)
class GoogleBaseTTSService(TTSService): class GoogleBaseTTSService(TTSService):
@@ -996,9 +994,7 @@ class GoogleTTSService(GoogleBaseTTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
error_message = f"TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message)
class GeminiTTSService(GoogleBaseTTSService): class GeminiTTSService(GoogleBaseTTSService):
@@ -1248,6 +1244,4 @@ class GeminiTTSService(GoogleBaseTTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
error_message = f"Gemini TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message)

View File

@@ -146,7 +146,6 @@ class GroqTTSService(TTSService):
bytes = w.readframes(num_frames) bytes = w.readframes(num_frames)
yield TTSAudioRawFrame(bytes, frame_rate, channels) yield TTSAudioRawFrame(bytes, frame_rate, channels)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -161,7 +161,7 @@ class HeyGenClient:
f"{self}::event_callback_task", f"{self}::event_callback_task",
) )
except Exception as e: except Exception as e:
logger.error(f"Failed to setup HeyGenClient: {e}") await self.push_error(error_msg=f"Failed to setup HeyGenClient: {e}", exception=e)
await self.cleanup() await self.cleanup()
async def cleanup(self) -> None: async def cleanup(self) -> None:
@@ -179,7 +179,7 @@ class HeyGenClient:
await self._task_manager.cancel_task(self._event_task) await self._task_manager.cancel_task(self._event_task)
self._event_task = None self._event_task = None
except Exception as e: except Exception as e:
logger.error(f"Exception during cleanup: {e}") await self.push_error(error_msg=f"Exception during cleanup: {e}", exception=e)
async def start(self, frame: StartFrame, audio_chunk_size: int) -> None: async def start(self, frame: StartFrame, audio_chunk_size: int) -> None:
"""Start the client and establish all necessary connections. """Start the client and establish all necessary connections.
@@ -229,7 +229,7 @@ class HeyGenClient:
self._ws_receive_task_handler(), name="HeyGenClient_Websocket" self._ws_receive_task_handler(), name="HeyGenClient_Websocket"
) )
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(error_msg=f"{self} initialization error: {e}", exception=e)
self._websocket = None self._websocket = None
async def _ws_receive_task_handler(self): async def _ws_receive_task_handler(self):
@@ -242,7 +242,9 @@ class HeyGenClient:
except ConnectionClosedOK: except ConnectionClosedOK:
break break
except Exception as e: except Exception as e:
logger.error(f"Error processing WebSocket message: {e}") await self.push_error(
error_msg=f"Error processing WebSocket message: {e}", exception=e
)
break break
async def _handle_ws_server_event(self, event: dict) -> None: async def _handle_ws_server_event(self, event: dict) -> None:
@@ -260,7 +262,7 @@ class HeyGenClient:
if self._websocket: if self._websocket:
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} disconnect error: {e}") await self.push_error(error_msg=f"{self} disconnect error: {e}", exception=e)
finally: finally:
self._websocket = None self._websocket = None
@@ -273,7 +275,9 @@ class HeyGenClient:
if self._websocket: if self._websocket:
await self._websocket.send(json.dumps(message)) await self._websocket.send(json.dumps(message))
except Exception as e: except Exception as e:
logger.error(f"Error sending message to HeyGen websocket: {e}") await self.push_error(
error_msg=f"Error sending message to HeyGen websocket: {e}", exception=e
)
raise e raise e
async def interrupt(self, event_id: str) -> None: async def interrupt(self, event_id: str) -> None:
@@ -471,9 +475,11 @@ class HeyGenClient:
await self._audio_frame_callback(audio_frame) await self._audio_frame_callback(audio_frame)
except Exception as e: except Exception as e:
logger.error(f"Error processing audio frame: {e}") await self.push_error(
error_msg=f"Error processing audio frame: {e}", exception=e
)
except Exception as e: except Exception as e:
logger.error(f"Error processing audio frames: {e}") await self.push_error(error_msg=f"Error processing audio frames: {e}", exception=e)
finally: finally:
logger.debug(f"Audio frame processing ended.") logger.debug(f"Audio frame processing ended.")
@@ -500,9 +506,11 @@ class HeyGenClient:
if self._transport_ready and self._video_frame_callback: if self._transport_ready and self._video_frame_callback:
await self._video_frame_callback(image_frame) await self._video_frame_callback(image_frame)
except Exception as e: except Exception as e:
logger.error(f"Error processing individual video frame: {e}") await self.push_error(
error_msg=f"Error processing individual video frame: {e}", exception=e
)
except Exception as e: except Exception as e:
logger.error(f"Error processing video frames: {e}") await self.push_error(error_msg=f"Error processing video frames: {e}", exception=e)
finally: finally:
logger.debug(f"Video frame processing ended.") logger.debug(f"Video frame processing ended.")
@@ -595,7 +603,7 @@ class HeyGenClient:
) )
except Exception as e: except Exception as e:
logger.error(f"LiveKit initialization error: {e}") await self.push_error(error_msg=f"LiveKit initialization error: {e}", exception=e)
self._livekit_room = None self._livekit_room = None
async def _livekit_disconnect(self): async def _livekit_disconnect(self):
@@ -616,7 +624,7 @@ class HeyGenClient:
self._livekit_room = None self._livekit_room = None
logger.debug("Successfully disconnected from LiveKit room") logger.debug("Successfully disconnected from LiveKit room")
except Exception as e: except Exception as e:
logger.error(f"LiveKit disconnect error: {e}") await self.push_error(error_msg=f"LiveKit disconnect error: {e}", exception=e)
# #
# Queue callback handling # Queue callback handling

View File

@@ -230,8 +230,7 @@ class LmntTTSService(InterruptibleTTSService):
# await self._websocket.send(json.dumps({"eof": True})) # await self._websocket.send(json.dumps({"eof": True}))
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error disconnecting from LMNT: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._started = False self._started = False
self._websocket = None self._websocket = None
@@ -300,13 +299,11 @@ class LmntTTSService(InterruptibleTTSService):
await self._get_websocket().send(json.dumps({"flush": True})) await self._get_websocket().send(json.dumps({"flush": True}))
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -176,7 +176,6 @@ class MCPClient(BaseObject):
except Exception as e: except Exception as e:
error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}" error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}"
logger.error(error_msg) logger.error(error_msg)
logger.exception("Full exception details:")
await params.result_callback(error_msg) await params.result_callback(error_msg)
async def _stdio_list_tools(self) -> ToolsSchema: async def _stdio_list_tools(self) -> ToolsSchema:
@@ -207,7 +206,6 @@ class MCPClient(BaseObject):
except Exception as e: except Exception as e:
error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}" error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}"
logger.error(error_msg) logger.error(error_msg)
logger.exception("Full exception details:")
await params.result_callback(error_msg) await params.result_callback(error_msg)
async def _streamable_http_list_tools(self) -> ToolsSchema: async def _streamable_http_list_tools(self) -> ToolsSchema:
@@ -246,7 +244,6 @@ class MCPClient(BaseObject):
except Exception as e: except Exception as e:
error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}" error_msg = f"Error calling mcp tool {params.function_name}: {str(e)}"
logger.error(error_msg) logger.error(error_msg)
logger.exception("Full exception details:")
await params.result_callback(error_msg) await params.result_callback(error_msg)
async def _call_tool(self, session, function_name, arguments, result_callback): async def _call_tool(self, session, function_name, arguments, result_callback):
@@ -302,7 +299,6 @@ class MCPClient(BaseObject):
except Exception as e: except Exception as e:
logger.error(f"Failed to read tool '{tool_name}': {str(e)}") logger.error(f"Failed to read tool '{tool_name}': {str(e)}")
logger.exception("Full exception details:")
continue continue
logger.debug(f"Completed reading {len(tool_schemas)} tools") logger.debug(f"Completed reading {len(tool_schemas)} tools")

View File

@@ -253,8 +253,7 @@ class Mem0MemoryService(FrameProcessor):
# Otherwise, pass the enhanced context frame downstream # Otherwise, pass the enhanced context frame downstream
await self.push_frame(frame) await self.push_frame(frame)
except Exception as e: except Exception as e:
logger.error(f"Error processing with Mem0: {str(e)}") await self.push_error(exception=e)
await self.push_frame(ErrorFrame(f"Error processing with Mem0: {str(e)}"))
await self.push_frame(frame) # Still pass the original frame through await self.push_frame(frame) # Still pass the original frame through
else: else:
# For non-context frames, just pass them through # For non-context frames, just pass them through

View File

@@ -264,8 +264,7 @@ class MiniMaxHttpTTSService(TTSService):
) as response: ) as response:
if response.status != 200: if response.status != 200:
error_message = f"MiniMax TTS error: HTTP {response.status}" error_message = f"MiniMax TTS error: HTTP {response.status}"
logger.error(error_message) await self.push_error(error_msg=error_message)
yield ErrorFrame(error=error_message)
return return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
@@ -338,8 +337,7 @@ class MiniMaxHttpTTSService(TTSService):
continue continue
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -110,8 +110,7 @@ class MoondreamService(VisionService):
if analysis fails. if analysis fails.
""" """
if not self._model: if not self._model:
logger.error(f"{self} error: Moondream model not available ({self.model_name})") await self.push_error(error_msg="Moondream model not initialized")
yield ErrorFrame("Moondream model not available")
return return
logger.debug(f"Analyzing image (bytes length: {len(frame.image)})") logger.debug(f"Analyzing image (bytes length: {len(frame.image)})")

View File

@@ -298,8 +298,7 @@ class NeuphonicTTSService(InterruptibleTTSService):
logger.debug("Disconnecting from Neuphonic") logger.debug("Disconnecting from Neuphonic")
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._started = False self._started = False
self._websocket = None self._websocket = None
@@ -364,16 +363,14 @@ class NeuphonicTTSService(InterruptibleTTSService):
await self._send_text(text) await self._send_text(text)
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class NeuphonicHttpTTSService(TTSService): class NeuphonicHttpTTSService(TTSService):
@@ -538,7 +535,6 @@ class NeuphonicHttpTTSService(TTSService):
error_text = await response.text() error_text = await response.text()
error_message = f"Neuphonic API error: HTTP {response.status} - {error_text}" error_message = f"Neuphonic API error: HTTP {response.status} - {error_text}"
logger.error(error_message) logger.error(error_message)
yield ErrorFrame(error=error_message)
return return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
@@ -567,8 +563,7 @@ class NeuphonicHttpTTSService(TTSService):
yield TTSAudioRawFrame(audio_bytes, self.sample_rate, 1) yield TTSAudioRawFrame(audio_bytes, self.sample_rate, 1)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
# Don't yield error frame for individual message failures # Don't yield error frame for individual message failures
continue continue
@@ -576,8 +571,7 @@ class NeuphonicHttpTTSService(TTSService):
logger.debug("TTS generation cancelled") logger.debug("TTS generation cancelled")
raise raise
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -76,8 +76,7 @@ class OpenAIImageGenService(ImageGenService):
image_url = image.data[0].url image_url = image.data[0].url
if not image_url: if not image_url:
logger.error(f"{self} No image provided in response: {image}") await self.push_error(error_msg="Image generation failed")
yield ErrorFrame("Image generation failed")
return return
# Load the image from the url # Load the image from the url

View File

@@ -443,7 +443,7 @@ class OpenAIRealtimeLLMService(LLMService):
) )
self._receive_task = self.create_task(self._receive_task_handler()) self._receive_task = self.create_task(self._receive_task_handler())
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
self._websocket = None self._websocket = None
async def _disconnect(self): async def _disconnect(self):
@@ -460,7 +460,7 @@ class OpenAIRealtimeLLMService(LLMService):
self._completed_tool_calls = set() self._completed_tool_calls = set()
self._disconnecting = False self._disconnecting = False
except Exception as e: except Exception as e:
logger.error(f"{self} error disconnecting: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
async def _ws_send(self, realtime_message): async def _ws_send(self, realtime_message):
try: try:
@@ -473,7 +473,6 @@ class OpenAIRealtimeLLMService(LLMService):
# somehow *started* the websocket send attempt while we still # somehow *started* the websocket send attempt while we still
# had a connection) # had a connection)
return return
logger.error(f"Error sending message to websocket: {e}")
# In server-to-server contexts, a WebSocket error should be quite rare. Given how hard # In server-to-server contexts, a WebSocket error should be quite rare. Given how hard
# it is to recover from a send-side error with proper state management, and that exponential # it is to recover from a send-side error with proper state management, and that exponential
# backoff for retries can have cost/stability implications for a service cluster, let's just # backoff for retries can have cost/stability implications for a service cluster, let's just

View File

@@ -206,5 +206,4 @@ class OpenAITTSService(TTSService):
yield frame yield frame
yield TTSStoppedFrame() yield TTSStoppedFrame()
except BadRequestError as e: except BadRequestError as e:
logger.error(f"{self} error generating TTS: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -79,5 +79,5 @@ class AzureRealtimeBetaLLMService(OpenAIRealtimeBetaLLMService):
) )
self._receive_task = self.create_task(self._receive_task_handler()) self._receive_task = self.create_task(self._receive_task_handler())
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
self._websocket = None self._websocket = None

View File

@@ -424,7 +424,7 @@ class OpenAIRealtimeBetaLLMService(LLMService):
) )
self._receive_task = self.create_task(self._receive_task_handler()) self._receive_task = self.create_task(self._receive_task_handler())
except Exception as e: except Exception as e:
logger.error(f"{self} initialization error: {e}") await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
self._websocket = None self._websocket = None
async def _disconnect(self): async def _disconnect(self):
@@ -440,7 +440,7 @@ class OpenAIRealtimeBetaLLMService(LLMService):
self._receive_task = None self._receive_task = None
self._disconnecting = False self._disconnecting = False
except Exception as e: except Exception as e:
logger.error(f"{self} error disconnecting: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
async def _ws_send(self, realtime_message): async def _ws_send(self, realtime_message):
try: try:
@@ -449,7 +449,6 @@ class OpenAIRealtimeBetaLLMService(LLMService):
except Exception as e: except Exception as e:
if self._disconnecting: if self._disconnecting:
return return
logger.error(f"Error sending message to websocket: {e}")
# In server-to-server contexts, a WebSocket error should be quite rare. Given how hard # In server-to-server contexts, a WebSocket error should be quite rare. Given how hard
# it is to recover from a send-side error with proper state management, and that exponential # it is to recover from a send-side error with proper state management, and that exponential
# backoff for retries can have cost/stability implications for a service cluster, let's just # backoff for retries can have cost/stability implications for a service cluster, let's just

View File

@@ -88,11 +88,8 @@ class PiperTTSService(TTSService):
) as response: ) as response:
if response.status != 200: if response.status != 200:
error = await response.text() error = await response.text()
logger.error( await self.push_error(
f"{self} error getting audio (status: {response.status}, error: {error})" error_msg=f"Error getting audio (status: {response.status}, error: {error})"
)
yield ErrorFrame(
error=f"Error getting audio (status: {response.status}, error: {error})"
) )
return return

View File

@@ -266,7 +266,7 @@ class PlayHTTTSService(InterruptibleTTSService):
self._websocket = None self._websocket = None
await self._call_event_handler("on_connection_error", f"{e}") await self._call_event_handler("on_connection_error", f"{e}")
except Exception as e: except Exception as e:
await self.push_error(exception=e) await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
self._websocket = None self._websocket = None
await self._call_event_handler("on_connection_error", f"{e}") await self._call_event_handler("on_connection_error", f"{e}")
@@ -279,8 +279,7 @@ class PlayHTTTSService(InterruptibleTTSService):
logger.debug("Disconnecting from PlayHT") logger.debug("Disconnecting from PlayHT")
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._request_id = None self._request_id = None
self._websocket = None self._websocket = None
@@ -350,7 +349,6 @@ class PlayHTTTSService(InterruptibleTTSService):
await self.push_frame(TTSStoppedFrame()) await self.push_frame(TTSStoppedFrame())
self._request_id = None self._request_id = None
elif "error" in msg: elif "error" in msg:
logger.error(f"{self} error: {msg}")
await self.push_error(error_msg=f"{self} error: {msg['error']}") await self.push_error(error_msg=f"{self} error: {msg['error']}")
except json.JSONDecodeError: except json.JSONDecodeError:
logger.error(f"Invalid JSON message: {message}") logger.error(f"Invalid JSON message: {message}")
@@ -393,8 +391,7 @@ class PlayHTTTSService(InterruptibleTTSService):
await self._get_websocket().send(json.dumps(tts_command)) await self._get_websocket().send(json.dumps(tts_command))
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
@@ -404,8 +401,7 @@ class PlayHTTTSService(InterruptibleTTSService):
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class PlayHTHttpTTSService(TTSService): class PlayHTHttpTTSService(TTSService):
@@ -625,8 +621,7 @@ class PlayHTHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -259,7 +259,7 @@ class RimeTTSService(AudioContextWordTTSService):
await self._call_event_handler("on_connected") await self._call_event_handler("on_connected")
except Exception as e: except Exception as e:
await self.push_error(exception=e) await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
self._websocket = None self._websocket = None
await self._call_event_handler("on_connection_error", f"{e}") await self._call_event_handler("on_connection_error", f"{e}")
@@ -271,8 +271,7 @@ class RimeTTSService(AudioContextWordTTSService):
await self._websocket.send(json.dumps(self._build_eos_msg())) await self._websocket.send(json.dumps(self._build_eos_msg()))
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error disconnecting: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._context_id = None self._context_id = None
self._websocket = None self._websocket = None
@@ -409,16 +408,14 @@ class RimeTTSService(AudioContextWordTTSService):
await self._get_websocket().send(json.dumps(msg)) await self._get_websocket().send(json.dumps(msg))
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class RimeHttpTTSService(TTSService): class RimeHttpTTSService(TTSService):
@@ -549,8 +546,7 @@ class RimeHttpTTSService(TTSService):
) as response: ) as response:
if response.status != 200: if response.status != 200:
error_message = f"Rime TTS error: HTTP {response.status}" error_message = f"Rime TTS error: HTTP {response.status}"
logger.error(error_message) await self.push_error(error_msg=error_message)
yield ErrorFrame(error=error_message)
return return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
@@ -567,8 +563,7 @@ class RimeHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -655,12 +655,12 @@ class RivaSegmentedSTTService(SegmentedSTTService):
logger.debug("No transcription results found in Riva response") logger.debug("No transcription results found in Riva response")
except AttributeError as ae: except AttributeError as ae:
logger.error(f"Unexpected response structure from Riva: {ae}") await self.push_error(
yield ErrorFrame(f"Unexpected Riva response format: {str(ae)}") error_msg=f"Unexpected response structure from Riva: {ae}", exception=ae
)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating STT: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
class ParakeetSTTService(RivaSTTService): class ParakeetSTTService(RivaSTTService):

View File

@@ -180,8 +180,13 @@ class RivaTTSService(TTSService):
yield frame yield frame
resp = await asyncio.wait_for(queue.get(), timeout=RIVA_TTS_TIMEOUT_SECS) resp = await asyncio.wait_for(queue.get(), timeout=RIVA_TTS_TIMEOUT_SECS)
except asyncio.TimeoutError: except asyncio.TimeoutError:
logger.error(f"{self} timeout waiting for audio response") await self.push_error(error_msg=f"Timeout generating TTS: {text}")
yield ErrorFrame(error=f"{self} error: {e}") yield TTSStoppedFrame()
return
except Exception as e:
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield TTSStoppedFrame()
return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -348,7 +348,9 @@ class SarvamSTTService(STTService):
# Exit the async context manager # Exit the async context manager
await self._websocket_context.__aexit__(None, None, None) await self._websocket_context.__aexit__(None, None, None)
except Exception as e: except Exception as e:
logger.error(f"Error closing WebSocket connection: {e}") await self.push_error(
error_msg=f"Error closing WebSocket connection: {e}", exception=e
)
finally: finally:
logger.debug("Disconnected from Sarvam WebSocket") logger.debug("Disconnected from Sarvam WebSocket")
self._socket_client = None self._socket_client = None
@@ -368,8 +370,7 @@ class SarvamSTTService(STTService):
# Messages will be handled via the _message_handler callback # Messages will be handled via the _message_handler callback
await self._socket_client.start_listening() await self._socket_client.start_listening()
except Exception as e: except Exception as e:
logger.error(f"Error in Sarvam receive task: {e}") await self.push_error(error_msg=f"Sarvam receive task error: {e}", exception=e)
await self.push_error(ErrorFrame(f"Sarvam receive task error: {e}"))
async def _handle_message(self, message): async def _handle_message(self, message):
"""Handle incoming WebSocket message from Sarvam SDK. """Handle incoming WebSocket message from Sarvam SDK.

View File

@@ -284,8 +284,7 @@ class SarvamHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()
yield TTSStoppedFrame() yield TTSStoppedFrame()
@@ -582,8 +581,9 @@ class SarvamTTSService(InterruptibleTTSService):
await self._call_event_handler("on_connected") await self._call_event_handler("on_connected")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(
await self.push_error(ErrorFrame(error=f"{self} error: {e}")) error_msg=f"Error connecting to Sarvam TTS Websocket: {e}", exception=e
)
self._websocket = None self._websocket = None
await self._call_event_handler("on_connection_error", f"{e}") await self._call_event_handler("on_connection_error", f"{e}")
@@ -611,8 +611,7 @@ class SarvamTTSService(InterruptibleTTSService):
logger.debug("Disconnecting from Sarvam") logger.debug("Disconnecting from Sarvam")
await self._websocket.close() await self._websocket.close()
except Exception as e: except Exception as e:
logger.error(f"{self} error closing websocket: {e}") await self.push_error(error_msg=f"Error closing websocket: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._started = False self._started = False
self._websocket = None self._websocket = None
@@ -636,7 +635,7 @@ class SarvamTTSService(InterruptibleTTSService):
await self.push_frame(frame) await self.push_frame(frame)
elif msg.get("type") == "error": elif msg.get("type") == "error":
error_msg = msg["data"]["message"] error_msg = msg["data"]["message"]
logger.error(f"TTS Error: {error_msg}") await self.push_error(error_msg=f"TTS Error: {error_msg}")
# If it's a timeout error, the connection might need to be reset # If it's a timeout error, the connection might need to be reset
if "too long" in error_msg.lower() or "timeout" in error_msg.lower(): if "too long" in error_msg.lower() or "timeout" in error_msg.lower():
@@ -698,13 +697,11 @@ class SarvamTTSService(InterruptibleTTSService):
await self._send_text(text) await self._send_text(text)
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error sending text: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
return return
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -178,7 +178,7 @@ class SimliVideoService(FrameProcessor):
self._audio_task = self.create_task(self._consume_and_process_audio()) self._audio_task = self.create_task(self._consume_and_process_audio())
self._video_task = self.create_task(self._consume_and_process_video()) self._video_task = self.create_task(self._consume_and_process_video())
except Exception as e: except Exception as e:
logger.error(f"{self}: unable to start connection: {e}") await self.push_error(error_msg=f"Unable to start connection: {e}", exception=e)
async def _consume_and_process_audio(self): async def _consume_and_process_audio(self):
"""Consume audio frames from Simli and push them downstream.""" """Consume audio frames from Simli and push them downstream."""
@@ -256,7 +256,7 @@ class SimliVideoService(FrameProcessor):
await self._simli_client.send(audioBytes) await self._simli_client.send(audioBytes)
return return
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error sending audio: {e}", exception=e)
elif isinstance(frame, TTSStoppedFrame): elif isinstance(frame, TTSStoppedFrame):
try: try:
if self._previously_interrupted and len(self._audio_buffer) > 0: if self._previously_interrupted and len(self._audio_buffer) > 0:
@@ -264,7 +264,7 @@ class SimliVideoService(FrameProcessor):
self._previously_interrupted = False self._previously_interrupted = False
self._audio_buffer = bytearray() self._audio_buffer = bytearray()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error stopping TTS: {e}", exception=e)
return return
elif isinstance(frame, (EndFrame, CancelFrame)): elif isinstance(frame, (EndFrame, CancelFrame)):
await self._stop() await self._stop()

View File

@@ -194,7 +194,7 @@ class SonioxSTTService(STTService):
self._websocket = await websocket_connect(self._url) self._websocket = await websocket_connect(self._url)
if not self._websocket: if not self._websocket:
logger.error(f"Unable to connect to Soniox API at {self._url}") await self.push_error(error_msg=f"Unable to connect to Soniox API at {self._url}")
# If vad_force_turn_endpoint is not enabled, we need to enable endpoint detection. # If vad_force_turn_endpoint is not enabled, we need to enable endpoint detection.
# Either one or the other is required. # Either one or the other is required.
@@ -419,5 +419,4 @@ class SonioxSTTService(STTService):
# Expected when closing the connection. # Expected when closing the connection.
pass pass
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error receiving message: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))

View File

@@ -467,8 +467,7 @@ class SpeechmaticsSTTService(STTService):
await self._client.send_audio(audio) await self._client.send_audio(audio)
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error sending audio: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
await self._disconnect() await self._disconnect()
def update_params( def update_params(
@@ -580,8 +579,7 @@ class SpeechmaticsSTTService(STTService):
logger.debug(f"{self} Connected to Speechmatics STT service") logger.debug(f"{self} Connected to Speechmatics STT service")
await self._call_event_handler("on_connected") await self._call_event_handler("on_connected")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error connecting to Speechmatics: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
self._client = None self._client = None
async def _disconnect(self) -> None: async def _disconnect(self) -> None:
@@ -595,7 +593,9 @@ class SpeechmaticsSTTService(STTService):
except asyncio.TimeoutError: except asyncio.TimeoutError:
logger.warning(f"{self} Timeout while closing Speechmatics client connection") logger.warning(f"{self} Timeout while closing Speechmatics client connection")
except Exception as e: except Exception as e:
await self.push_error(exception=e) await self.push_error(
error_msg=f"Error disconnecting from Speechmatics: {e}", exception=e
)
finally: finally:
self._client = None self._client = None
await self._call_event_handler("on_disconnected") await self._call_event_handler("on_disconnected")

View File

@@ -436,18 +436,12 @@ class UltravoxSTTService(AIService):
yield LLMFullResponseEndFrame() yield LLMFullResponseEndFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error generating text: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
else: else:
logger.error("No model available for text generation") await self.push_error(error_msg="No model available for text generation")
yield ErrorFrame("No model available for text generation")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error processing audio: {e}", exception=e)
import traceback
logger.error(traceback.format_exc())
yield ErrorFrame(f"Error processing audio: {str(e)}")
finally: finally:
self._buffer.is_processing = False self._buffer.is_processing = False
self._buffer.frames = [] self._buffer.frames = []

View File

@@ -226,8 +226,7 @@ class BaseWhisperSTTService(SegmentedSTTService):
logger.warning("Received empty transcription from API") logger.warning("Received empty transcription from API")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(exception=e)
yield ErrorFrame(error=f"{self} error: {e}")
async def _transcribe(self, audio: bytes) -> Transcription: async def _transcribe(self, audio: bytes) -> Transcription:
"""Transcribe audio data to text. """Transcribe audio data to text.

View File

@@ -285,8 +285,7 @@ class WhisperSTTService(SegmentedSTTService):
The service will normalize it to float32 in the range [-1, 1]. The service will normalize it to float32 in the range [-1, 1].
""" """
if not self._model: if not self._model:
logger.error(f"{self} error: Whisper model not available") await self.push_error(error_msg="Whisper model not available")
yield ErrorFrame("Whisper model not available")
return return
await self.start_processing_metrics() await self.start_processing_metrics()
@@ -428,5 +427,4 @@ class WhisperSTTServiceMLX(WhisperSTTService):
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error processing audio: {e}", exception=e)
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -181,8 +181,9 @@ class XTTSService(TTSService):
async with self._aiohttp_session.post(url, json=payload) as r: async with self._aiohttp_session.post(url, json=payload) as r:
if r.status != 200: if r.status != 200:
text = await r.text() text = await r.text()
logger.error(f"{self} error getting audio (status: {r.status}, error: {text})") await self.push_error(
yield ErrorFrame(error=f"Error getting audio (status: {r.status}, error: {text})") error_msg=f"Error getting audio (status: {r.status}, error: {text})"
)
return return
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)