Enhanced error handling across the framework.

Co-authored-by: Mark Backman <m.backman@gmail.com>
This commit is contained in:
Filipi Fuchter
2025-11-26 18:34:25 -03:00
parent 9efb21d61e
commit 1330ef3ad6
74 changed files with 268 additions and 373 deletions

View File

@@ -9,6 +9,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Added ### Added
- Enhanced error handling across the framework:
- Added `on_error` callback to `FrameProcessor` for centralized error handling.
- Renamed `push_error(error: ErrorFrame)` to `push_error_frame(error: ErrorFrame)`
for clarity.
- Added new `push_error` method for simplified error reporting:
```
push_error(error_msg: str, exception: Optional[Exception] = None, fatal: bool = False)
```
- Standardized error logging by replacing `logger.exception` calls with
`logger.error` throughout the codebase.
- Added `cache_read_input_tokens`, `cache_creation_input_tokens` and - Added `cache_read_input_tokens`, `cache_creation_input_tokens` and
`reasoning_tokens` to OTel spans for LLM call `reasoning_tokens` to OTel spans for LLM call

View File

@@ -835,11 +835,13 @@ class ErrorFrame(SystemFrame):
error: Description of the error that occurred. error: Description of the error that occurred.
fatal: Whether the error is fatal and requires bot shutdown. fatal: Whether the error is fatal and requires bot shutdown.
processor: The frame processor that generated the error. processor: The frame processor that generated the error.
exception: The exception that occurred.
""" """
error: str error: str
fatal: bool = False fatal: bool = False
processor: Optional["FrameProcessor"] = None processor: Optional["FrameProcessor"] = None
exception: Optional[Exception] = None
def __str__(self): def __str__(self):
return f"{self.name}(error: {self.error}, fatal: {self.fatal})" return f"{self.name}(error: {self.error}, fatal: {self.fatal})"

View File

@@ -126,6 +126,4 @@ class WakeCheckFilter(FrameProcessor):
else: else:
await self.push_frame(frame, direction) await self.push_frame(frame, direction)
except Exception as e: except Exception as e:
error_msg = f"Error in wake word filter: {e}" await self.push_error(error_msg=f"Error in wake word filter: {e}", exception=e)
logger.exception(error_msg)
await self.push_error(ErrorFrame(error_msg))

View File

@@ -142,6 +142,7 @@ class FrameProcessor(BaseObject):
- on_after_process_frame: Called after a frame is processed - on_after_process_frame: Called after a frame is processed
- on_before_push_frame: Called before a frame is pushed - on_before_push_frame: Called before a frame is pushed
- on_after_push_frame: Called after a frame is pushed - on_after_push_frame: Called after a frame is pushed
- on_error: Called when an error is raised in the frame processing.
""" """
def __init__( def __init__(
@@ -234,6 +235,7 @@ class FrameProcessor(BaseObject):
self._register_event_handler("on_after_process_frame", sync=True) self._register_event_handler("on_after_process_frame", sync=True)
self._register_event_handler("on_before_push_frame", sync=True) self._register_event_handler("on_before_push_frame", sync=True)
self._register_event_handler("on_after_push_frame", sync=True) self._register_event_handler("on_after_push_frame", sync=True)
self._register_event_handler("on_error", sync=True)
@property @property
def id(self) -> int: def id(self) -> int:
@@ -630,7 +632,43 @@ class FrameProcessor(BaseObject):
elif isinstance(frame, (FrameProcessorResumeFrame, FrameProcessorResumeUrgentFrame)): elif isinstance(frame, (FrameProcessorResumeFrame, FrameProcessorResumeUrgentFrame)):
await self.__resume(frame) await self.__resume(frame)
async def push_error(self, error: ErrorFrame): async def push_error(
self,
error_msg: str,
exception: Optional[Exception] = None,
fatal: bool = False,
):
"""Creates and pushes an ErrorFrame upstream.
Creates and pushes an ErrorFrame upstream to notify other processors in the
pipeline about an error condition. The error frame will include context about
which processor generated the error.
Args:
error_msg: Descriptive message explaining the error condition.
exception: Optional exception object that caused the error, if available.
This provides additional context for debugging and error handling.
fatal: Whether this error should be considered fatal to the pipeline.
Fatal errors typically cause the entire pipeline to stop processing.
Defaults to False for non-fatal errors.
Example::
```python
# Non-fatal error
await self.push_error("Failed to process audio chunk, skipping")
# Fatal error with exception context
try:
result = some_critical_operation()
except Exception as e:
await self.push_error("Critical operation failed", exception=e, fatal=True)
```
"""
error_frame = ErrorFrame(error=error_msg, fatal=fatal, exception=exception, processor=self)
await self.push_error_frame(error=error_frame)
async def push_error_frame(self, error: ErrorFrame):
"""Push an error frame upstream. """Push an error frame upstream.
Args: Args:
@@ -638,6 +676,8 @@ class FrameProcessor(BaseObject):
""" """
if not error.processor: if not error.processor:
error.processor = self error.processor = self
await self._call_event_handler("on_error", error)
logger.error(f"{error.processor} error: {error.error}")
await self.push_frame(error, FrameDirection.UPSTREAM) await self.push_frame(error, FrameDirection.UPSTREAM)
async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM): async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM):
@@ -759,8 +799,10 @@ class FrameProcessor(BaseObject):
await self.__cancel_process_task() await self.__cancel_process_task()
self.__create_process_task() self.__create_process_task()
except Exception as e: except Exception as e:
logger.exception(f"Uncaught exception in {self} when handling _start_interruption: {e}") await self.push_error(
await self.push_error(ErrorFrame(str(e))) error_msg=f"Uncaught exception handling _start_interruption: {e}",
exception=e,
)
async def __internal_push_frame(self, frame: Frame, direction: FrameDirection): async def __internal_push_frame(self, frame: Frame, direction: FrameDirection):
"""Internal method to push frames to adjacent processors. """Internal method to push frames to adjacent processors.
@@ -797,8 +839,7 @@ class FrameProcessor(BaseObject):
await self._observer.on_push_frame(data) await self._observer.on_push_frame(data)
await self._prev.queue_frame(frame, direction) await self._prev.queue_frame(frame, direction)
except Exception as e: except Exception as e:
logger.exception(f"Uncaught exception in {self}: {e}") await self.push_error(error_msg=f"Uncaught exception: {e}", exception=e)
await self.push_error(ErrorFrame(str(e)))
def _check_started(self, frame: Frame): def _check_started(self, frame: Frame):
"""Check if the processor has been started. """Check if the processor has been started.
@@ -874,8 +915,7 @@ class FrameProcessor(BaseObject):
await self._call_event_handler("on_after_process_frame", frame) await self._call_event_handler("on_after_process_frame", frame)
except Exception as e: except Exception as e:
logger.exception(f"{self}: error processing frame: {e}") await self.push_error(error_msg=f"Error processing frame: {e}", exception=e)
await self.push_error(ErrorFrame(str(e)))
async def __input_frame_task_handler(self): async def __input_frame_task_handler(self):
"""Handle frames from the input queue. """Handle frames from the input queue.

View File

@@ -24,7 +24,7 @@ try:
from langchain_core.messages import AIMessageChunk from langchain_core.messages import AIMessageChunk
from langchain_core.runnables import Runnable from langchain_core.runnables import Runnable
except ModuleNotFoundError as e: except ModuleNotFoundError as e:
logger.exception("In order to use Langchain, you need to `pip install pipecat-ai[langchain]`. ") logger.error("In order to use Langchain, you need to `pip install pipecat-ai[langchain]`. ")
raise Exception(f"Missing module: {e}") raise Exception(f"Missing module: {e}")
@@ -113,6 +113,6 @@ class LangchainProcessor(FrameProcessor):
except GeneratorExit: except GeneratorExit:
logger.warning(f"{self} generator was closed prematurely") logger.warning(f"{self} generator was closed prematurely")
except Exception as e: except Exception as e:
logger.exception(f"{self} an unknown error occurred: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
finally: finally:
await self.push_frame(LLMFullResponseEndFrame()) await self.push_frame(LLMFullResponseEndFrame())

View File

@@ -23,7 +23,7 @@ try:
from strands import Agent from strands import Agent
from strands.multiagent.graph import Graph from strands.multiagent.graph import Graph
except ModuleNotFoundError as e: except ModuleNotFoundError as e:
logger.exception("In order to use Strands Agents, you need to `pip install strands-agents`.") logger.error("In order to use Strands Agents, you need to `pip install strands-agents`.")
raise Exception(f"Missing module: {e}") raise Exception(f"Missing module: {e}")
@@ -143,7 +143,7 @@ class StrandsAgentsProcessor(FrameProcessor):
except GeneratorExit: except GeneratorExit:
logger.warning(f"{self} generator was closed prematurely") logger.warning(f"{self} generator was closed prematurely")
except Exception as e: except Exception as e:
logger.exception(f"{self} an unknown error occurred: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
finally: finally:
if ttfb_tracking: if ttfb_tracking:
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()

View File

@@ -199,7 +199,7 @@ class PlivoFrameSerializer(FrameSerializer):
) )
except Exception as e: except Exception as e:
logger.exception(f"Failed to hang up Plivo call: {e}") logger.error(f"Failed to hang up Plivo call: {e}")
async def deserialize(self, data: str | bytes) -> Frame | None: async def deserialize(self, data: str | bytes) -> Frame | None:
"""Deserializes Plivo WebSocket data to Pipecat frames. """Deserializes Plivo WebSocket data to Pipecat frames.

View File

@@ -225,7 +225,7 @@ class TelnyxFrameSerializer(FrameSerializer):
) )
except Exception as e: except Exception as e:
logger.exception(f"Failed to hang up Telnyx call: {e}") logger.error(f"Failed to hang up Telnyx call: {e}")
async def deserialize(self, data: str | bytes) -> Frame | None: async def deserialize(self, data: str | bytes) -> Frame | None:
"""Deserializes Telnyx WebSocket data to Pipecat frames. """Deserializes Telnyx WebSocket data to Pipecat frames.

View File

@@ -236,7 +236,7 @@ class TwilioFrameSerializer(FrameSerializer):
) )
except Exception as e: except Exception as e:
logger.exception(f"Failed to hang up Twilio call: {e}") logger.error(f"Failed to hang up Twilio call: {e}")
async def deserialize(self, data: str | bytes) -> Frame | None: async def deserialize(self, data: str | bytes) -> Frame | None:
"""Deserializes Twilio WebSocket data to Pipecat frames. """Deserializes Twilio WebSocket data to Pipecat frames.

View File

@@ -166,6 +166,6 @@ class AIService(FrameProcessor):
async for f in generator: async for f in generator:
if f: if f:
if isinstance(f, ErrorFrame): if isinstance(f, ErrorFrame):
await self.push_error(f) await self.push_error_frame(f)
else: else:
await self.push_frame(f) await self.push_frame(f)

View File

@@ -458,8 +458,7 @@ class AnthropicLLMService(LLMService):
except httpx.TimeoutException: except httpx.TimeoutException:
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.exception(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(f"{e}"))
finally: finally:
await self.stop_processing_metrics() await self.stop_processing_metrics()
await self.push_frame(LLMFullResponseEndFrame()) await self.push_frame(LLMFullResponseEndFrame())

View File

@@ -206,9 +206,8 @@ class AssemblyAISTTService(STTService):
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._connected = False self._connected = False
await self.push_error(ErrorFrame(error=f"{self} error: {e}")) await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
raise raise
async def _disconnect(self): async def _disconnect(self):
@@ -233,8 +232,7 @@ class AssemblyAISTTService(STTService):
logger.warning("Timed out waiting for termination message from server") logger.warning("Timed out waiting for termination message from server")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
if self._receive_task: if self._receive_task:
await self.cancel_task(self._receive_task) await self.cancel_task(self._receive_task)
@@ -242,8 +240,7 @@ class AssemblyAISTTService(STTService):
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._websocket = None self._websocket = None
@@ -262,13 +259,11 @@ class AssemblyAISTTService(STTService):
except websockets.exceptions.ConnectionClosedOK: except websockets.exceptions.ConnectionClosedOK:
break break
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
break break
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
def _parse_message(self, message: Dict[str, Any]) -> BaseMessage: def _parse_message(self, message: Dict[str, Any]) -> BaseMessage:
"""Parse a raw message into the appropriate message type.""" """Parse a raw message into the appropriate message type."""
@@ -297,8 +292,7 @@ class AssemblyAISTTService(STTService):
elif isinstance(parsed_message, TerminationMessage): elif isinstance(parsed_message, TerminationMessage):
await self._handle_termination(parsed_message) await self._handle_termination(parsed_message)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def _handle_termination(self, message: TerminationMessage): async def _handle_termination(self, message: TerminationMessage):
"""Handle termination message.""" """Handle termination message."""

View File

@@ -228,8 +228,7 @@ class AsyncAITTSService(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(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -241,8 +240,7 @@ class AsyncAITTSService(InterruptibleTTSService):
logger.debug("Disconnecting from Async") logger.debug("Disconnecting from Async")
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._websocket = None self._websocket = None
self._started = False self._started = False
@@ -287,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(ErrorFrame(error=f"{self} error: {msg['message']}")) await self.push_error(error_msg=f"Error: {msg['message']}")
else: else:
logger.error(f"{self} error, unknown message type: {msg}") await self.push_error(error_msg=f"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."""
@@ -335,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class AsyncAIHttpTTSService(TTSService): class AsyncAIHttpTTSService(TTSService):
@@ -477,8 +472,7 @@ class AsyncAIHttpTTSService(TTSService):
async with self._session.post(url, json=payload, headers=headers) as response: async with self._session.post(url, json=payload, headers=headers) as response:
if response.status != 200: if response.status != 200:
error_text = await response.text() error_text = await response.text()
logger.error(f"Async API error: {error_text}") await self.push_error(error_msg=f"Async API error: {error_text}")
await self.push_error(ErrorFrame(error=f"Async API error: {error_text}"))
raise Exception(f"Async API returned status {response.status}: {error_text}") raise Exception(f"Async API returned status {response.status}: {error_text}")
audio_data = await response.read() audio_data = await response.read()
@@ -494,8 +488,7 @@ class AsyncAIHttpTTSService(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"Unknown error occurred: {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()

View File

@@ -1136,7 +1136,7 @@ class AWSBedrockLLMService(LLMService):
except (ReadTimeoutError, asyncio.TimeoutError): except (ReadTimeoutError, asyncio.TimeoutError):
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.exception(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
finally: finally:
await self.stop_processing_metrics() await self.stop_processing_metrics()
await self.push_frame(LLMFullResponseEndFrame()) await self.push_frame(LLMFullResponseEndFrame())

View File

@@ -453,7 +453,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(error_msg=f"Initialization error: {e}", 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):
@@ -577,7 +577,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(
@@ -885,7 +885,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(error_msg=f"Error processing responses: {e}", exception=e)
if self._wants_connection: if self._wants_connection:
await self.reset_conversation() await self.reset_conversation()

View File

@@ -140,8 +140,7 @@ class AWSTranscribeSTTService(STTService):
return return
logger.warning("WebSocket connection not established after connect") logger.warning("WebSocket connection not established after connect")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
retry_count += 1 retry_count += 1
if retry_count < max_retries: if retry_count < max_retries:
await asyncio.sleep(1) # Wait before retrying await asyncio.sleep(1) # Wait before retrying
@@ -182,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}") yield ErrorFrame(error=f"Unknown error occurred: {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
@@ -200,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
await self._disconnect() await self._disconnect()
async def _connect(self): async def _connect(self):
@@ -289,8 +285,7 @@ class AWSTranscribeSTTService(STTService):
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
await self._disconnect() await self._disconnect()
raise raise
@@ -310,8 +305,7 @@ class AWSTranscribeSTTService(STTService):
await self._ws_client.send(json.dumps(end_stream)) await self._ws_client.send(json.dumps(end_stream))
await self._ws_client.close() await self._ws_client.close()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._ws_client = None self._ws_client = None
await self._call_event_handler("on_disconnected") await self._call_event_handler("on_disconnected")
@@ -529,15 +523,15 @@ 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:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
break break

View File

@@ -312,7 +312,6 @@ class AWSPollyTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except (BotoCoreError, ClientError) as error: except (BotoCoreError, ClientError) as error:
logger.exception(f"{self} error generating TTS: {error}")
error_message = f"AWS Polly TTS error: {str(error)}" error_message = f"AWS Polly TTS error: {str(error)}"
yield ErrorFrame(error=error_message) yield ErrorFrame(error=error_message)

View File

@@ -91,7 +91,6 @@ 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")
yield ErrorFrame("Image generation timed out") yield ErrorFrame("Image generation timed out")
return return
@@ -104,7 +103,6 @@ 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")
yield ErrorFrame("Image generation failed") yield ErrorFrame("Image generation failed")
return return

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(error_msg=f"initialization error: {e}", 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}") yield ErrorFrame(error=f"Unknown error occurred: {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.
@@ -151,8 +150,9 @@ class AzureSTTService(STTService):
self._speech_recognizer.recognized.connect(self._on_handle_recognized) self._speech_recognizer.recognized.connect(self._on_handle_recognized)
self._speech_recognizer.start_continuous_recognition_async() self._speech_recognizer.start_continuous_recognition_async()
except Exception as e: except Exception as e:
logger.error(f"{self} exception during initialization: {e}") await self.push_error(
await self.push_error(ErrorFrame(error=f"{self} error: {e}")) error_msg=f"Uncaught exception during initialization: {e}", exception=e
)
async def stop(self, frame: EndFrame): async def stop(self, frame: EndFrame):
"""Stop the speech recognition service. """Stop the speech recognition service.

View File

@@ -327,7 +327,6 @@ 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)
yield ErrorFrame(error=error_msg) yield ErrorFrame(error=error_msg)
return return
@@ -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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class AzureHttpTTSService(AzureBaseTTSService): class AzureHttpTTSService(AzureBaseTTSService):
@@ -440,5 +437,6 @@ 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}") yield ErrorFrame(
yield ErrorFrame(error=f"{self} error: {cancellation_details.error_details}") error=f"Unknown error occurred: {cancellation_details.error_details}"
)

View File

@@ -276,8 +276,7 @@ class CartesiaSTTService(WebsocketSTTService):
self._websocket = await websocket_connect(ws_url, additional_headers=headers) self._websocket = await websocket_connect(ws_url, additional_headers=headers)
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def _disconnect_websocket(self): async def _disconnect_websocket(self):
try: try:
@@ -285,8 +284,7 @@ class CartesiaSTTService(WebsocketSTTService):
logger.debug("Disconnecting from Cartesia STT") logger.debug("Disconnecting from Cartesia 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)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._websocket = None self._websocket = None
await self._call_event_handler("on_disconnected") await self._call_event_handler("on_disconnected")
@@ -319,8 +317,7 @@ class CartesiaSTTService(WebsocketSTTService):
elif data["type"] == "error": elif data["type"] == "error":
error_msg = data.get("message", "Unknown error") error_msg = data.get("message", "Unknown error")
logger.error(f"Cartesia error: {error_msg}") await self.push_error(error_msg=error_msg)
await self.push_error(ErrorFrame(error=error_msg))
@traced_stt @traced_stt
async def _handle_transcription( async def _handle_transcription(

View File

@@ -497,8 +497,7 @@ class CartesiaTTSService(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}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -510,8 +509,7 @@ class CartesiaTTSService(AudioContextWordTTSService):
logger.debug("Disconnecting from Cartesia") logger.debug("Disconnecting from Cartesia")
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"Unknown error occurred: {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
@@ -564,13 +562,12 @@ class CartesiaTTSService(AudioContextWordTTSService):
) )
await self.append_to_audio_context(msg["context_id"], frame) await self.append_to_audio_context(msg["context_id"], frame)
elif msg["type"] == "error": elif msg["type"] == "error":
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(ErrorFrame(error=f"{self} error: {msg['error']}")) await self.push_error(error_msg=f"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"Error, unknown message type: {msg}")
async def _receive_messages(self): async def _receive_messages(self):
while True: while True:
@@ -608,16 +605,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class CartesiaHttpTTSService(TTSService): class CartesiaHttpTTSService(TTSService):
@@ -808,8 +803,7 @@ class CartesiaHttpTTSService(TTSService):
async with session.post(url, json=payload, headers=headers) as response: async with session.post(url, json=payload, headers=headers) as response:
if response.status != 200: if response.status != 200:
error_text = await response.text() error_text = await response.text()
logger.error(f"Cartesia API error: {error_text}") yield ErrorFrame(error=f"Cartesia API error: {error_text}")
await self.push_error(ErrorFrame(error=f"Cartesia API error: {error_text}"))
raise Exception(f"Cartesia API returned status {response.status}: {error_text}") raise Exception(f"Cartesia API returned status {response.status}: {error_text}")
audio_data = await response.read() audio_data = await response.read()
@@ -825,8 +819,7 @@ class CartesiaHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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()

View File

@@ -192,8 +192,7 @@ class DeepgramFluxSTTService(WebsocketSTTService):
try: try:
await self._disconnect_websocket() await self._disconnect_websocket()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
# Reset state only after everything is cleaned up # Reset state only after everything is cleaned up
self._websocket = None self._websocket = None
@@ -251,8 +250,7 @@ class DeepgramFluxSTTService(WebsocketSTTService):
logger.debug("Connected to Deepgram Flux Websocket") logger.debug("Connected to Deepgram Flux Websocket")
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -280,8 +278,7 @@ class DeepgramFluxSTTService(WebsocketSTTService):
logger.debug("Disconnecting from Deepgram Flux Websocket") logger.debug("Disconnecting from Deepgram Flux Websocket")
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._websocket = None self._websocket = None
await self._call_event_handler("on_disconnected") await self._call_event_handler("on_disconnected")
@@ -381,7 +378,6 @@ 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.")
yield ErrorFrame("Not connected to Deepgram Flux.") yield ErrorFrame("Not connected to Deepgram Flux.")
return return
@@ -389,8 +385,7 @@ class DeepgramFluxSTTService(WebsocketSTTService):
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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
return return
yield None yield None
@@ -467,8 +462,7 @@ class DeepgramFluxSTTService(WebsocketSTTService):
# Skip malformed messages # Skip malformed messages
continue continue
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
# Error will be handled inside WebsocketService->_receive_task_handler # Error will be handled inside WebsocketService->_receive_task_handler
raise raise
else: else:

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():
@@ -256,7 +256,7 @@ class DeepgramSTTService(STTService):
async def _on_error(self, *args, **kwargs): async def _on_error(self, *args, **kwargs):
error: ErrorResponse = kwargs["error"] error: ErrorResponse = kwargs["error"]
logger.warning(f"{self} connection error, will retry: {error}") logger.warning(f"{self} connection error, will retry: {error}")
await self.push_error(ErrorFrame(error=f"{error}")) await self.push_error(error_msg=f"{error}")
await self.stop_all_metrics() await self.stop_all_metrics()
# NOTE(aleix): we don't disconnect (i.e. call finish on the connection) # NOTE(aleix): we don't disconnect (i.e. call finish on the connection)
# because this triggers more errors internally in the Deepgram SDK. So, # because this triggers more errors internally in the Deepgram SDK. So,

View File

@@ -290,8 +290,7 @@ class DeepgramTTSService(WebsocketTTSService):
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class DeepgramHttpTTSService(TTSService): class DeepgramHttpTTSService(TTSService):
@@ -401,5 +400,4 @@ class DeepgramHttpTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.exception(f"{self} exception: {e}")
yield ErrorFrame(f"Error getting audio: {str(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}") yield ErrorFrame(error=f"Unknown error occurred: {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:
@@ -598,7 +597,6 @@ 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}")
yield ErrorFrame(f"ElevenLabs Realtime STT error: {str(e)}") yield ErrorFrame(f"ElevenLabs Realtime STT error: {str(e)}")
yield None yield None
@@ -663,8 +661,9 @@ class ElevenLabsRealtimeSTTService(WebsocketSTTService):
await self._call_event_handler("on_connected") await self._call_event_handler("on_connected")
logger.debug("Connected to ElevenLabs Realtime STT") logger.debug("Connected to ElevenLabs Realtime STT")
except Exception as e: except Exception as e:
logger.error(f"{self}: unable to connect to ElevenLabs Realtime STT: {e}") await self.push_error(
await self.push_error(ErrorFrame(f"Connection error: {str(e)}")) error_msg=f"Unable to connect to ElevenLabs Realtime STT: {e}", exception=e
)
async def _disconnect_websocket(self): async def _disconnect_websocket(self):
"""Disconnect from ElevenLabs Realtime STT WebSocket.""" """Disconnect from ElevenLabs Realtime STT WebSocket."""
@@ -673,7 +672,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")
@@ -733,17 +732,17 @@ class ElevenLabsRealtimeSTTService(WebsocketSTTService):
elif message_type == "error": elif message_type == "error":
error_msg = data.get("error", "Unknown error") error_msg = data.get("error", "Unknown error")
logger.error(f"ElevenLabs error: {error_msg}") logger.error(f"ElevenLabs error: {error_msg}")
await self.push_error(ErrorFrame(f"Error: {error_msg}")) await self.push_error(error_msg=f"Error: {error_msg}")
elif message_type == "auth_error": elif message_type == "auth_error":
error_msg = data.get("error", "Authentication error") error_msg = data.get("error", "Authentication error")
logger.error(f"ElevenLabs auth error: {error_msg}") logger.error(f"ElevenLabs auth error: {error_msg}")
await self.push_error(ErrorFrame(f"Auth error: {error_msg}")) await self.push_error(error_msg=f"Auth error: {error_msg}")
elif message_type == "quota_exceeded_error": elif message_type == "quota_exceeded_error":
error_msg = data.get("error", "Quota exceeded") error_msg = data.get("error", "Quota exceeded")
logger.error(f"ElevenLabs quota exceeded: {error_msg}") logger.error(f"ElevenLabs quota exceeded: {error_msg}")
await self.push_error(ErrorFrame(f"Quota exceeded: {error_msg}")) await self.push_error(error_msg=f"Quota exceeded: {error_msg}")
else: else:
logger.debug(f"Unknown message type: {message_type}") logger.debug(f"Unknown message type: {message_type}")

View File

@@ -424,8 +424,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(error_msg=f"Unknown error occurred: {e}", 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
@@ -536,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(error_msg=f"Unknown error occurred: {e}", 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):
@@ -553,8 +551,7 @@ class ElevenLabsTTSService(AudioContextWordTTSService):
await self._websocket.close() await self._websocket.close()
logger.debug("Disconnected from ElevenLabs") logger.debug("Disconnected from ElevenLabs")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._started = False self._started = False
self._context_id = None self._context_id = None
@@ -584,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(error_msg=f"Unknown error occurred: {e}", 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 = ""
@@ -740,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}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
yield ErrorFrame(error=f"{self} error: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class ElevenLabsHttpTTSService(WordTTSService): class ElevenLabsHttpTTSService(WordTTSService):
@@ -1043,7 +1037,6 @@ 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}")
yield ErrorFrame(error=f"ElevenLabs API error: {error_text}") yield ErrorFrame(error=f"ElevenLabs API error: {error_text}")
return return
@@ -1091,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}") yield ErrorFrame(error=f"Unknown error occurred: {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
@@ -1116,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}") yield ErrorFrame(error=f"Unknown error occurred: {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,7 +110,6 @@ 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")
yield ErrorFrame("Image generation failed") yield ErrorFrame("Image generation failed")
return return

View File

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

View File

@@ -228,8 +228,7 @@ class FishAudioTTSService(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(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -243,8 +242,7 @@ class FishAudioTTSService(InterruptibleTTSService):
await self._websocket.send(ormsgpack.packb(stop_message)) await self._websocket.send(ormsgpack.packb(stop_message))
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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
self._request_id = None self._request_id = None
self._started = False self._started = False
@@ -286,8 +284,7 @@ class FishAudioTTSService(InterruptibleTTSService):
continue continue
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
@traced_tts @traced_tts
async def run_tts(self, text: str) -> AsyncGenerator[Frame, None]: async def run_tts(self, text: str) -> AsyncGenerator[Frame, None]:
@@ -323,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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
@@ -332,5 +328,4 @@ class FishAudioTTSService(InterruptibleTTSService):
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -468,8 +468,7 @@ class GladiaSTTService(STTService):
break break
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
self._connection_active = False self._connection_active = False
if not self._should_reconnect: if not self._should_reconnect:
@@ -559,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(error_msg=f"Unknown error occurred: {e}", 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:
@@ -623,8 +621,7 @@ class GladiaSTTService(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"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def _maybe_reconnect(self) -> bool: async def _maybe_reconnect(self) -> bool:
"""Handle exponential backoff reconnection logic.""" """Handle exponential backoff reconnection logic."""
@@ -632,7 +629,9 @@ 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",
)
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

@@ -1175,7 +1175,7 @@ class GeminiLiveLLMService(LLMService):
self._connection_task = self.create_task(self._connection_task_handler(config=config)) self._connection_task = self.create_task(self._connection_task_handler(config=config))
except Exception as e: except Exception as e:
await self.push_error(ErrorFrame(error=f"{self} Initialization error: {e}")) await self.push_error(error_msg=f"Initialization error: {e}", exception=e)
async def _connection_task_handler(self, config: LiveConnectConfig): async def _connection_task_handler(self, config: LiveConnectConfig):
async with self._client.aio.live.connect(model=self._model_name, config=config) as session: async with self._client.aio.live.connect(model=self._model_name, config=config) as session:
@@ -1252,11 +1252,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)
return False return False
else: else:
logger.info( logger.info(
@@ -1284,7 +1284,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."""
@@ -1745,7 +1745,7 @@ class GeminiLiveLLMService(LLMService):
# state management, and that exponential backoff for retries can have # state management, and that exponential backoff for retries can have
# cost/stability implications for a service cluster, let's just treat a # cost/stability implications for a service cluster, let's just treat a
# send-side error as fatal. # send-side error as fatal.
await self.push_error(ErrorFrame(error=f"{self} Send error: {error}", fatal=True)) await self.push_error(error_msg=f"Send error: {error}")
def create_context_aggregator( def create_context_aggregator(
self, self,

View File

@@ -110,7 +110,6 @@ 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")
yield ErrorFrame("Image generation failed") yield ErrorFrame("Image generation failed")
return return
@@ -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}")
yield ErrorFrame(f"Image generation error: {str(e)}") yield ErrorFrame(f"Image generation error: {str(e)}")

View File

@@ -793,7 +793,7 @@ class GoogleLLMService(LLMService):
return return
generation_params.setdefault("thinking_config", {})["thinking_budget"] = 0 generation_params.setdefault("thinking_config", {})["thinking_budget"] = 0
except Exception as e: except Exception as e:
logger.exception(f"Failed to unset thinking budget: {e}") logger.error(f"Failed to unset thinking budget: {e}")
async def _stream_content( async def _stream_content(
self, params_from_context: GeminiLLMInvocationParams self, params_from_context: GeminiLLMInvocationParams
@@ -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.exception(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", 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

@@ -774,8 +774,7 @@ class GoogleSTTService(STTService):
yield cloud_speech.StreamingRecognizeRequest(audio=audio_data) yield cloud_speech.StreamingRecognizeRequest(audio=audio_data)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
raise raise
async def _stream_audio(self): async def _stream_audio(self):
@@ -806,15 +805,13 @@ class GoogleSTTService(STTService):
break break
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", 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)
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def run_stt(self, audio: bytes) -> AsyncGenerator[Frame, None]: async def run_stt(self, audio: bytes) -> AsyncGenerator[Frame, None]:
"""Process an audio chunk for STT transcription. """Process an audio chunk for STT transcription.
@@ -902,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(error_msg=f"Unknown error occurred: {e}", 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,7 +737,6 @@ class GoogleHttpTTSService(TTSService):
yield TTSStoppedFrame() yield TTSStoppedFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}")
error_message = f"TTS generation error: {str(e)}" error_message = f"TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message) yield ErrorFrame(error=error_message)
@@ -996,9 +995,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(error_msg=f"TTS generation error: {str(e)}", exception=e)
error_message = f"TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message)
class GeminiTTSService(GoogleBaseTTSService): class GeminiTTSService(GoogleBaseTTSService):
@@ -1248,6 +1245,5 @@ class GeminiTTSService(GoogleBaseTTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}")
error_message = f"Gemini TTS generation error: {str(e)}" error_message = f"Gemini TTS generation error: {str(e)}"
yield ErrorFrame(error=error_message) 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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -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.exception(f"Exception during cleanup: {e}") logger.error(f"Exception during cleanup: {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.

View File

@@ -287,8 +287,7 @@ class HumeTTSService(WordTTSService):
self._cumulative_time = utterance_duration self._cumulative_time = utterance_duration
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
# Ensure TTFB timer is stopped even on early failures # Ensure TTFB timer is stopped even on early failures
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()

View File

@@ -397,8 +397,7 @@ class InworldTTSService(TTSService):
# STEP 7: ERROR HANDLING # STEP 7: ERROR HANDLING
# ================================================================================ # ================================================================================
# Log any unexpected errors and notify the pipeline # Log any unexpected errors and notify the pipeline
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
# ================================================================================ # ================================================================================
# STEP 8: CLEANUP AND COMPLETION # STEP 8: CLEANUP AND COMPLETION
@@ -513,7 +512,7 @@ class InworldTTSService(TTSService):
# Extract the base64-encoded audio content from response # Extract the base64-encoded audio content from response
if "audioContent" not in response_data: if "audioContent" not in response_data:
logger.error("No audioContent in Inworld API response") logger.error("No audioContent in Inworld API response")
await self.push_error(ErrorFrame(error="No audioContent in response")) yield ErrorFrame(error="No audioContent in response")
return return
# ================================================================================ # ================================================================================

View File

@@ -214,8 +214,7 @@ class LmntTTSService(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(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -231,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
@@ -266,10 +264,9 @@ class LmntTTSService(InterruptibleTTSService):
try: try:
msg = json.loads(message) msg = json.loads(message)
if "error" in msg: if "error" in msg:
logger.error(f"{self} error: {msg['error']}")
await self.push_frame(TTSStoppedFrame()) await self.push_frame(TTSStoppedFrame())
await self.stop_all_metrics() await self.stop_all_metrics()
await self.push_error(ErrorFrame(error=f"{self} error: {msg['error']}")) await self.push_error(error_msg=f"Error: {msg['error']}")
return return
except json.JSONDecodeError: except json.JSONDecodeError:
logger.error(f"Invalid JSON message: {message}") logger.error(f"Invalid JSON message: {message}")
@@ -302,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {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,9 @@ 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(
await self.push_frame(ErrorFrame(f"Error processing with Mem0: {str(e)}")) error_msg=f"Error processing with Mem0: {str(e)}", exception=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

@@ -314,7 +314,6 @@ 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)
yield ErrorFrame(error=error_message) yield ErrorFrame(error=error_message)
return return
@@ -392,8 +391,7 @@ class MiniMaxHttpTTSService(TTSService):
continue continue
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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

@@ -110,7 +110,6 @@ 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})")
yield ErrorFrame("Moondream model not available") yield ErrorFrame("Moondream model not available")
return return

View File

@@ -285,8 +285,7 @@ class NeuphonicTTSService(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(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -299,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(error_msg=f"Unknown error occurred: {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
@@ -365,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class NeuphonicHttpTTSService(TTSService): class NeuphonicHttpTTSService(TTSService):
@@ -538,7 +534,6 @@ class NeuphonicHttpTTSService(TTSService):
if response.status != 200: if response.status != 200:
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)
yield ErrorFrame(error=error_message) yield ErrorFrame(error=error_message)
return return
@@ -568,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}") yield ErrorFrame(error=f"Unknown error occurred: {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
@@ -577,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}") yield ErrorFrame(error=f"Unknown error occurred: {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,7 +76,6 @@ 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}")
yield ErrorFrame("Image generation failed") yield ErrorFrame("Image generation failed")
return return

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,12 +473,11 @@ 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
# treat a send-side error as fatal. # treat a send-side error as fatal.
await self.push_error(ErrorFrame(error=f"Error sending client event: {e}")) await self.push_error(error_msg=f"Error sending client event: {e}", exception=e)
async def _update_settings(self): async def _update_settings(self):
settings = self._session_properties settings = self._session_properties
@@ -674,7 +673,7 @@ class OpenAIRealtimeLLMService(LLMService):
self._current_assistant_response = None self._current_assistant_response = None
# error handling # error handling
if evt.response.status == "failed": if evt.response.status == "failed":
await self.push_error(ErrorFrame(error=evt.response.status_details["error"]["message"])) await self.push_error(error_msg=evt.response.status_details["error"]["message"])
return return
# response content # response content
for item in evt.response.output: for item in evt.response.output:
@@ -766,7 +765,7 @@ class OpenAIRealtimeLLMService(LLMService):
async def _handle_evt_error(self, evt): async def _handle_evt_error(self, evt):
# Errors are fatal to this connection. Send an ErrorFrame. # Errors are fatal to this connection. Send an ErrorFrame.
await self.push_error(ErrorFrame(error=f"Error: {evt}")) await self.push_error(error_msg=f"Error: {evt}")
# #
# state and client events for the current conversation # state and client events for the current conversation

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.exception(f"{self} error generating TTS: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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

@@ -425,7 +425,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):
@@ -441,7 +441,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:
@@ -450,12 +450,11 @@ 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
# treat a send-side error as fatal. # treat a send-side error as fatal.
await self.push_error(ErrorFrame(error=f"Error sending client event: {e}")) await self.push_error(error_msg=f"Error sending client event: {e}", exception=e)
async def _update_settings(self): async def _update_settings(self):
settings = self._session_properties settings = self._session_properties
@@ -686,7 +685,7 @@ class OpenAIRealtimeBetaLLMService(LLMService):
async def _handle_evt_error(self, evt): async def _handle_evt_error(self, evt):
# Errors are fatal to this connection. Send an ErrorFrame. # Errors are fatal to this connection. Send an ErrorFrame.
await self.push_error(ErrorFrame(error=f"Error: {evt}")) await self.push_error(error_msg=f"Error: {evt}")
async def _handle_assistant_output(self, output): async def _handle_assistant_output(self, output):
# We haven't seen intermixed audio and function_call items in the same response. But let's # We haven't seen intermixed audio and function_call items in the same response. But let's

View File

@@ -88,9 +88,6 @@ 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(
f"{self} error getting audio (status: {response.status}, error: {error})"
)
yield ErrorFrame( yield ErrorFrame(
error=f"Error getting audio (status: {response.status}, error: {error})" error=f"Error getting audio (status: {response.status}, error: {error})"
) )
@@ -109,7 +106,7 @@ class PiperTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") logger.error(f"{self} exception: {e}")
yield ErrorFrame(error=f"{self} error: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
finally: finally:
logger.debug(f"{self}: Finished TTS [{text}]") logger.debug(f"{self}: Finished TTS [{text}]")
await self.stop_ttfb_metrics() await self.stop_ttfb_metrics()

View File

@@ -266,8 +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:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -280,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
@@ -351,8 +349,7 @@ 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"Error: {msg['error']}")
await self.push_error(ErrorFrame(error=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}")
@@ -394,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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
yield TTSStoppedFrame() yield TTSStoppedFrame()
await self._disconnect() await self._disconnect()
await self._connect() await self._connect()
@@ -405,8 +401,7 @@ class PlayHTTTSService(InterruptibleTTSService):
yield None yield None
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class PlayHTHttpTTSService(TTSService): class PlayHTHttpTTSService(TTSService):
@@ -626,8 +621,7 @@ class PlayHTHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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

@@ -300,8 +300,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:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Error connecting: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {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}")
@@ -313,8 +312,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
@@ -407,10 +405,9 @@ class RimeTTSService(AudioContextWordTTSService):
logger.debug(f"Updated cumulative time to: {self._cumulative_time}") logger.debug(f"Updated cumulative time to: {self._cumulative_time}")
elif msg["type"] == "error": elif msg["type"] == "error":
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(ErrorFrame(error=f"{self} error: {msg['message']}")) await self.push_error(error_msg=f"Error: {msg['message']}")
self._context_id = None self._context_id = None
async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM): async def push_frame(self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM):
@@ -452,16 +449,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class RimeHttpTTSService(TTSService): class RimeHttpTTSService(TTSService):
@@ -592,7 +587,6 @@ 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)
yield ErrorFrame(error=error_message) yield ErrorFrame(error=error_message)
return return
@@ -610,8 +604,7 @@ class RimeHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {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,10 @@ 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}")
yield ErrorFrame(f"Unexpected Riva response format: {str(ae)}") yield ErrorFrame(f"Unexpected Riva response format: {str(ae)}")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
class ParakeetSTTService(RivaSTTService): class ParakeetSTTService(RivaSTTService):

View File

@@ -180,8 +180,7 @@ 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") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
await self.start_tts_usage_metrics(text) await self.start_tts_usage_metrics(text)
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -275,8 +275,7 @@ class SarvamSTTService(STTService):
await self._socket_client.translate(**method_kwargs) await self._socket_client.translate(**method_kwargs)
except Exception as e: except Exception as e:
logger.error(f"Error sending audio to Sarvam: {e}") yield ErrorFrame(error=f"Error sending audio to Sarvam: {e}", exception=e)
await self.push_error(ErrorFrame(f"Failed to send audio: {e}"))
yield None yield None
@@ -332,13 +331,11 @@ class SarvamSTTService(STTService):
logger.info("Connected to Sarvam successfully") logger.info("Connected to Sarvam successfully")
except ApiError as e: except ApiError as e:
logger.error(f"Sarvam API error: {e}") await self.push_error(error_msg=f"Sarvam API error: {e}", exception=e)
await self.push_error(ErrorFrame(f"Sarvam API error: {e}"))
except Exception as e: except Exception as e:
logger.error(f"Failed to connect to Sarvam: {e}")
self._socket_client = None self._socket_client = None
self._websocket_context = None self._websocket_context = None
await self.push_error(ErrorFrame(f"Failed to connect to Sarvam: {e}")) await self.push_error(error_msg=f"Failed to connect to Sarvam: {e}", exception=e)
async def _disconnect(self): async def _disconnect(self):
"""Disconnect from Sarvam WebSocket API using SDK.""" """Disconnect from Sarvam WebSocket API using SDK."""
@@ -351,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
@@ -371,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.
@@ -427,8 +425,7 @@ class SarvamSTTService(STTService):
await self.stop_processing_metrics() await self.stop_processing_metrics()
except Exception as e: except Exception as e:
logger.error(f"Error handling Sarvam message: {e}") await self.push_error(error_msg=f"Failed to handle message: {e}", exception=e)
await self.push_error(ErrorFrame(f"Failed to handle message: {e}"))
await self.stop_all_metrics() await self.stop_all_metrics()
@traced_stt @traced_stt

View File

@@ -254,8 +254,7 @@ class SarvamHttpTTSService(TTSService):
async with self._session.post(url, json=payload, headers=headers) as response: async with self._session.post(url, json=payload, headers=headers) as response:
if response.status != 200: if response.status != 200:
error_text = await response.text() error_text = await response.text()
logger.error(f"Sarvam API error: {error_text}") yield ErrorFrame(error=f"Sarvam API error: {error_text}")
await self.push_error(ErrorFrame(error=f"Sarvam API error: {error_text}"))
return return
response_data = await response.json() response_data = await response.json()
@@ -264,8 +263,7 @@ class SarvamHttpTTSService(TTSService):
# Decode base64 audio data # Decode base64 audio data
if "audios" not in response_data or not response_data["audios"]: if "audios" not in response_data or not response_data["audios"]:
logger.error("No audio data received from Sarvam API") yield ErrorFrame(error="No audio data received")
await self.push_error(ErrorFrame(error="No audio data received"))
return return
# Get the first audio (there should be only one for single text input) # Get the first audio (there should be only one for single text input)
@@ -286,8 +284,7 @@ class SarvamHttpTTSService(TTSService):
yield frame yield frame
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=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()
@@ -560,8 +557,7 @@ class SarvamTTSService(InterruptibleTTSService):
await self._disconnect_websocket() await self._disconnect_websocket()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
finally: finally:
# Reset state only after everything is cleaned up # Reset state only after everything is cleaned up
self._started = False self._started = False
@@ -585,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}")
@@ -602,8 +599,7 @@ class SarvamTTSService(InterruptibleTTSService):
await self._websocket.send(json.dumps(config_message)) await self._websocket.send(json.dumps(config_message))
logger.debug("Configuration sent successfully") logger.debug("Configuration sent successfully")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
raise raise
async def _disconnect_websocket(self): async def _disconnect_websocket(self):
@@ -615,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
@@ -640,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():
@@ -702,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}") yield ErrorFrame(error=f"Unknown error occurred: {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}") yield ErrorFrame(error=f"Unknown error occurred: {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.exception(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.exception(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.
@@ -327,8 +327,7 @@ class SonioxSTTService(STTService):
# Expected when closing the connection # Expected when closing the connection
logger.debug("WebSocket connection closed, keepalive task stopped.") logger.debug("WebSocket connection closed, keepalive task stopped.")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
async def _receive_task_handler(self): async def _receive_task_handler(self):
if not self._websocket: if not self._websocket:
@@ -404,13 +403,8 @@ class SonioxSTTService(STTService):
if error_code or error_message: if error_code or error_message:
# In case of error, still send the final transcript (if any remaining in the buffer). # In case of error, still send the final transcript (if any remaining in the buffer).
await send_endpoint_transcript() await send_endpoint_transcript()
logger.error(
f"{self} error: {error_code} (_receive_task_handler) - {error_message}"
)
await self.push_error( await self.push_error(
ErrorFrame( error_msg=f"Error: {error_code} (_receive_task_handler) - {error_message}"
error=f"{self} error: {error_code} (_receive_task_handler) - {error_message}"
)
) )
finished = content.get("finished") finished = content.get("finished")
@@ -425,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}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
await self._disconnect() await self._disconnect()
def update_params( def update_params(
@@ -514,8 +513,7 @@ class SpeechmaticsSTTService(STTService):
self._client.send_message(payload), self.get_event_loop() self._client.send_message(payload), self.get_event_loop()
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
raise RuntimeError(f"error sending message to STT: {e}") raise RuntimeError(f"error sending message to STT: {e}")
async def _connect(self) -> None: async def _connect(self) -> None:
@@ -581,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:
@@ -596,8 +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:
logger.error(f"{self} exception: {e}") await self.push_error(
await self.push_error(ErrorFrame(error=f"{self} error: {e}")) 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

@@ -163,7 +163,7 @@ class SpeechmaticsTTSService(TTSService):
# Report error frame # Report error frame
yield ErrorFrame( yield ErrorFrame(
error=f"{self} Service unavailable [503] (attempt {attempt}, retry in {backoff_time:.2f}s)" error=f"Service unavailable [503] (attempt {attempt}, retry in {backoff_time:.2f}s)"
) )
# Wait before retrying # Wait before retrying
@@ -174,16 +174,13 @@ class SpeechmaticsTTSService(TTSService):
except (ValueError, ArithmeticError): except (ValueError, ArithmeticError):
yield ErrorFrame( yield ErrorFrame(
error=f"{self} Service unavailable [503] (attempts {attempt})", error=f"Service unavailable [503] (attempts {attempt})",
fatal=True,
) )
return return
# != 200 : Error # != 200 : Error
if response.status != 200: if response.status != 200:
yield ErrorFrame( yield ErrorFrame(error=f"Service unavailable [{response.status}]")
error=f"{self} Service unavailable [{response.status}]", fatal=True
)
return return
# Update Pipecat metrics # Update Pipecat metrics
@@ -225,7 +222,7 @@ class SpeechmaticsTTSService(TTSService):
break break
except Exception as e: except Exception as e:
yield ErrorFrame(error=f"{self}: Error generating TTS: {e}", fatal=True) yield ErrorFrame(error=f"Error generating TTS: {e}")
finally: finally:
# Emit the TTS stopped frame # Emit the TTS stopped frame
yield TTSStoppedFrame() yield TTSStoppedFrame()

View File

@@ -329,4 +329,4 @@ class WebsocketSTTService(STTService, WebsocketService):
async def _report_error(self, error: ErrorFrame): async def _report_error(self, error: ErrorFrame):
await self._call_event_handler("on_connection_error", error.error) await self._call_event_handler("on_connection_error", error.error)
await self.push_error(error) await self.push_error_frame(error)

View File

@@ -781,7 +781,7 @@ class WebsocketTTSService(TTSService, WebsocketService):
async def _report_error(self, error: ErrorFrame): async def _report_error(self, error: ErrorFrame):
await self._call_event_handler("on_connection_error", error.error) await self._call_event_handler("on_connection_error", error.error)
await self.push_error(error) await self.push_error_frame(error)
class InterruptibleTTSService(WebsocketTTSService): class InterruptibleTTSService(WebsocketTTSService):
@@ -843,7 +843,7 @@ class WebsocketWordTTSService(WordTTSService, WebsocketService):
async def _report_error(self, error: ErrorFrame): async def _report_error(self, error: ErrorFrame):
await self._call_event_handler("on_connection_error", error.error) await self._call_event_handler("on_connection_error", error.error)
await self.push_error(error) await self.push_error_frame(error)
class InterruptibleWordTTSService(WebsocketWordTTSService): class InterruptibleWordTTSService(WebsocketWordTTSService):

View File

@@ -246,8 +246,7 @@ class UltravoxSTTService(AIService):
logger.info("Model warm-up completed successfully") logger.info("Model warm-up completed successfully")
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") await self.push_error(error_msg=f"Unknown error occurred: {e}", exception=e)
await self.push_error(ErrorFrame(error=f"{self} error: {e}"))
def _generate_silent_audio(self, sample_rate=16000, duration_sec=1.0): def _generate_silent_audio(self, sample_rate=16000, duration_sec=1.0):
"""Generate silent audio as a numpy array. """Generate silent audio as a numpy array.
@@ -377,7 +376,7 @@ class UltravoxSTTService(AIService):
if arr.size > 0: # Check if array is not empty if arr.size > 0: # Check if array is not empty
audio_arrays.append(arr) audio_arrays.append(arr)
except Exception as e: except Exception as e:
yield ErrorFrame(error=f"{self} error: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
# Handle numpy array data # Handle numpy array data
elif isinstance(f.audio, np.ndarray): elif isinstance(f.audio, np.ndarray):
if f.audio.size > 0: # Check if array is not empty if f.audio.size > 0: # Check if array is not empty
@@ -437,17 +436,11 @@ class UltravoxSTTService(AIService):
yield LLMFullResponseEndFrame() yield LLMFullResponseEndFrame()
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")
else: else:
logger.error("No model available for text generation")
yield ErrorFrame("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}")
import traceback
logger.error(traceback.format_exc())
yield ErrorFrame(f"Error processing audio: {str(e)}") yield ErrorFrame(f"Error processing audio: {str(e)}")
finally: finally:
self._buffer.is_processing = False self._buffer.is_processing = False

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}") yield ErrorFrame(error=f"Unknown error occurred: {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,7 +285,6 @@ 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")
yield ErrorFrame("Whisper model not available") yield ErrorFrame("Whisper model not available")
return return
@@ -428,5 +427,4 @@ class WhisperSTTServiceMLX(WhisperSTTService):
) )
except Exception as e: except Exception as e:
logger.error(f"{self} exception: {e}") yield ErrorFrame(error=f"Unknown error occurred: {e}")
yield ErrorFrame(error=f"{self} error: {e}")

View File

@@ -141,13 +141,8 @@ class XTTSService(TTSService):
async with self._aiohttp_session.get(self._settings["base_url"] + "/studio_speakers") as r: async with self._aiohttp_session.get(self._settings["base_url"] + "/studio_speakers") as r:
if r.status != 200: if r.status != 200:
text = await r.text() text = await r.text()
logger.error(
f"{self} error getting studio speakers (status: {r.status}, error: {text})"
)
await self.push_error( await self.push_error(
ErrorFrame( error_msg=f"Error getting studio speakers (status: {r.status}, error: {text})"
error=f"Error getting studio speakers (status: {r.status}, error: {text})"
)
) )
return return
self._studio_speakers = await r.json() self._studio_speakers = await r.json()
@@ -186,7 +181,6 @@ 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})")
yield ErrorFrame(error=f"Error getting audio (status: {r.status}, error: {text})") yield ErrorFrame(error=f"Error getting audio (status: {r.status}, error: {text})")
return return

View File

@@ -2506,13 +2506,10 @@ class DailyTransport(BaseTransport):
async def _on_error(self, error): async def _on_error(self, error):
"""Handle error events and push error frames.""" """Handle error events and push error frames."""
await self._call_event_handler("on_error", error) await self._call_event_handler("on_error", error)
# Push error frame to notify the pipeline
error_frame = ErrorFrame(error)
if self._input: if self._input:
await self._input.push_error(error_frame) await self._input.push_error(error_msg=error)
elif self._output: elif self._output:
await self._output.push_error(error_frame) await self._output.push_error(error_msg=error)
else: else:
logger.error("Both input and output are None while trying to push error") logger.error("Both input and output are None while trying to push error")
raise Exception("No valid input or output channel to push error") raise Exception("No valid input or output channel to push error")
@@ -2568,7 +2565,7 @@ class DailyTransport(BaseTransport):
except asyncio.TimeoutError: except asyncio.TimeoutError:
logger.error(f"Timeout handling dialin-ready event ({url})") logger.error(f"Timeout handling dialin-ready event ({url})")
except Exception as e: except Exception as e:
logger.exception(f"Error handling dialin-ready event ({url}): {e}") logger.error(f"Error handling dialin-ready event ({url}): {e}")
async def _on_dialin_connected(self, data): async def _on_dialin_connected(self, data):
"""Handle dial-in connected events.""" """Handle dial-in connected events."""

View File

@@ -316,7 +316,7 @@ class SmallWebRTCConnection(BaseObject):
logger.debug("Client not connected. Queuing app-message.") logger.debug("Client not connected. Queuing app-message.")
self._pending_app_messages.append(json_message) self._pending_app_messages.append(json_message)
except Exception as e: except Exception as e:
logger.exception(f"Error parsing JSON message {message}, {e}") logger.error(f"Error parsing JSON message {message}, {e}")
# Despite the fact that aiortc provides this listener, they don't have a status for "disconnected" # Despite the fact that aiortc provides this listener, they don't have a status for "disconnected"
# So, in case we loose connection, this event will not be triggered # So, in case we loose connection, this event will not be triggered

View File

@@ -265,7 +265,7 @@ class TavusTransportClient:
try: try:
await self._client.cleanup() await self._client.cleanup()
except Exception as e: except Exception as e:
logger.exception(f"Exception during cleanup: {e}") logger.error(f"Exception during cleanup: {e}")
async def _on_joined(self, data): async def _on_joined(self, data):
"""Handle joined event.""" """Handle joined event."""

View File

@@ -162,7 +162,7 @@ class TaskManager(BaseTaskManager):
# Re-raise the exception to ensure the task is cancelled. # Re-raise the exception to ensure the task is cancelled.
raise raise
except Exception as e: except Exception as e:
logger.exception(f"{name}: unexpected exception: {e}") logger.error(f"{name}: unexpected exception: {e}")
if not self._params: if not self._params:
raise Exception("TaskManager is not setup: unable to get event loop") raise Exception("TaskManager is not setup: unable to get event loop")
@@ -197,7 +197,7 @@ class TaskManager(BaseTaskManager):
# Here are sure the task is cancelled properly. # Here are sure the task is cancelled properly.
pass pass
except Exception as e: except Exception as e:
logger.exception(f"{name}: unexpected exception while cancelling task: {e}") logger.error(f"{name}: unexpected exception while cancelling task: {e}")
except BaseException as e: except BaseException as e:
logger.critical(f"{name}: fatal base exception while cancelling task: {e}") logger.critical(f"{name}: fatal base exception while cancelling task: {e}")
raise raise

View File

@@ -187,7 +187,7 @@ class BaseObject(ABC):
else: else:
handler(self, *args, **kwargs) handler(self, *args, **kwargs)
except Exception as e: except Exception as e:
logger.exception(f"Exception in event handler {event_name}: {e}") logger.error(f"Exception in event handler {event_name}: {e}")
def _event_task_finished(self, task: asyncio.Task): def _event_task_finished(self, task: asyncio.Task):
"""Clean up completed event handler tasks. """Clean up completed event handler tasks.