Restoring to use yield ErrorFrame
This commit is contained in:
@@ -181,7 +181,7 @@ class AWSTranscribeSTTService(STTService):
|
|||||||
try:
|
try:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
return
|
return
|
||||||
|
|
||||||
# Format the audio data according to AWS event stream format
|
# Format the audio data according to AWS event stream format
|
||||||
@@ -198,11 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
|
|
||||||
async def _connect(self):
|
async def _connect(self):
|
||||||
|
|||||||
@@ -311,8 +311,9 @@ class AWSPollyTTSService(TTSService):
|
|||||||
yield frame
|
yield frame
|
||||||
|
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
except (BotoCoreError, ClientError) as e:
|
except (BotoCoreError, ClientError) as error:
|
||||||
await self.push_error(exception=e)
|
error_message = f"AWS Polly TTS error: {str(error)}"
|
||||||
|
yield ErrorFrame(error=error_message)
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -91,7 +91,7 @@ class AzureImageGenServiceREST(ImageGenService):
|
|||||||
while status != "succeeded":
|
while status != "succeeded":
|
||||||
attempts_left -= 1
|
attempts_left -= 1
|
||||||
if attempts_left == 0:
|
if attempts_left == 0:
|
||||||
await self.push_error(error_msg="Image generation timed out")
|
yield ErrorFrame("Image generation timed out")
|
||||||
return
|
return
|
||||||
|
|
||||||
await asyncio.sleep(1)
|
await asyncio.sleep(1)
|
||||||
@@ -103,7 +103,7 @@ class AzureImageGenServiceREST(ImageGenService):
|
|||||||
|
|
||||||
image_url = json_response["result"]["data"][0]["url"] if json_response else None
|
image_url = json_response["result"]["data"][0]["url"] if json_response else None
|
||||||
if not image_url:
|
if not image_url:
|
||||||
await self.push_error(error_msg="Image generation failed")
|
yield ErrorFrame("Image generation failed")
|
||||||
return
|
return
|
||||||
|
|
||||||
# Load the image from the url
|
# Load the image from the url
|
||||||
|
|||||||
@@ -121,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
async def start(self, frame: StartFrame):
|
async def start(self, frame: StartFrame):
|
||||||
"""Start the speech recognition service.
|
"""Start the speech recognition service.
|
||||||
|
|||||||
@@ -327,7 +327,7 @@ class AzureTTSService(AzureBaseTTSService):
|
|||||||
try:
|
try:
|
||||||
if self._speech_synthesizer is None:
|
if self._speech_synthesizer is None:
|
||||||
error_msg = "Speech synthesizer not initialized."
|
error_msg = "Speech synthesizer not initialized."
|
||||||
await self.push_error(error_msg=error_msg)
|
yield ErrorFrame(error=error_msg)
|
||||||
return
|
return
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -354,13 +354,13 @@ class AzureTTSService(AzureBaseTTSService):
|
|||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
# Could add reconnection logic here if needed
|
# Could add reconnection logic here if needed
|
||||||
return
|
return
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class AzureHttpTTSService(AzureBaseTTSService):
|
class AzureHttpTTSService(AzureBaseTTSService):
|
||||||
@@ -437,4 +437,4 @@ class AzureHttpTTSService(AzureBaseTTSService):
|
|||||||
cancellation_details = result.cancellation_details
|
cancellation_details = result.cancellation_details
|
||||||
logger.warning(f"Speech synthesis canceled: {cancellation_details.reason}")
|
logger.warning(f"Speech synthesis canceled: {cancellation_details.reason}")
|
||||||
if cancellation_details.reason == CancellationReason.Error:
|
if cancellation_details.reason == CancellationReason.Error:
|
||||||
await self.push_error(error_msg=cancellation_details.error_details)
|
yield ErrorFrame(error=f"{self} error: {cancellation_details.error_details}")
|
||||||
|
|||||||
@@ -505,14 +505,14 @@ class CartesiaTTSService(AudioContextWordTTSService):
|
|||||||
await self._get_websocket().send(msg)
|
await self._get_websocket().send(msg)
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class CartesiaHttpTTSService(TTSService):
|
class CartesiaHttpTTSService(TTSService):
|
||||||
|
|||||||
@@ -378,14 +378,14 @@ class DeepgramFluxSTTService(WebsocketSTTService):
|
|||||||
are issues sending the audio data.
|
are issues sending the audio data.
|
||||||
"""
|
"""
|
||||||
if not self._websocket:
|
if not self._websocket:
|
||||||
await self.push_error(error_msg=f"Websocket not connected")
|
yield ErrorFrame("Not connected to Deepgram Flux.")
|
||||||
return
|
return
|
||||||
|
|
||||||
try:
|
try:
|
||||||
self._last_stt_time = time.monotonic()
|
self._last_stt_time = time.monotonic()
|
||||||
await self.send_with_retry(audio, self._report_error)
|
await self.send_with_retry(audio, self._report_error)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
return
|
return
|
||||||
|
|
||||||
yield None
|
yield None
|
||||||
|
|||||||
@@ -116,7 +116,7 @@ class DeepgramTTSService(TTSService):
|
|||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class DeepgramHttpTTSService(TTSService):
|
class DeepgramHttpTTSService(TTSService):
|
||||||
@@ -226,4 +226,4 @@ class DeepgramHttpTTSService(TTSService):
|
|||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(f"Error getting audio: {str(e)}")
|
||||||
|
|||||||
@@ -351,7 +351,7 @@ class ElevenLabsSTTService(SegmentedSTTService):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
def audio_format_from_sample_rate(sample_rate: int) -> str:
|
def audio_format_from_sample_rate(sample_rate: int) -> str:
|
||||||
@@ -585,7 +585,7 @@ class ElevenLabsRealtimeSTTService(WebsocketSTTService):
|
|||||||
}
|
}
|
||||||
await self._websocket.send(json.dumps(message))
|
await self._websocket.send(json.dumps(message))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(f"ElevenLabs Realtime STT error: {str(e)}")
|
||||||
|
|
||||||
yield None
|
yield None
|
||||||
|
|
||||||
|
|||||||
@@ -736,13 +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:
|
||||||
await self.push_error(exception=e)
|
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
self._started = False
|
self._started = False
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class ElevenLabsHttpTTSService(WordTTSService):
|
class ElevenLabsHttpTTSService(WordTTSService):
|
||||||
@@ -1037,7 +1037,7 @@ class ElevenLabsHttpTTSService(WordTTSService):
|
|||||||
) as response:
|
) as response:
|
||||||
if response.status != 200:
|
if response.status != 200:
|
||||||
error_text = await response.text()
|
error_text = await response.text()
|
||||||
await self.push_error(error_msg=f"ElevenLabs API error: {error_text}")
|
yield ErrorFrame(error=f"ElevenLabs API error: {error_text}")
|
||||||
return
|
return
|
||||||
|
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
@@ -1084,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
continue
|
continue
|
||||||
|
|
||||||
# After processing all chunks, emit any remaining partial word
|
# After processing all chunks, emit any remaining partial word
|
||||||
@@ -1108,7 +1108,7 @@ class ElevenLabsHttpTTSService(WordTTSService):
|
|||||||
self._previous_text = text
|
self._previous_text = text
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
finally:
|
finally:
|
||||||
await self.stop_ttfb_metrics()
|
await self.stop_ttfb_metrics()
|
||||||
# Let the parent class handle TTSStoppedFrame
|
# Let the parent class handle TTSStoppedFrame
|
||||||
|
|||||||
@@ -110,7 +110,7 @@ class FalImageGenService(ImageGenService):
|
|||||||
image_url = response["images"][0]["url"] if response else None
|
image_url = response["images"][0]["url"] if response else None
|
||||||
|
|
||||||
if not image_url:
|
if not image_url:
|
||||||
await self.push_error(error_msg="Image generation failed")
|
yield ErrorFrame("Image generation failed")
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug(f"Image generated at: {image_url}")
|
logger.debug(f"Image generated at: {image_url}")
|
||||||
|
|||||||
@@ -290,4 +290,4 @@ class FalSTTService(SegmentedSTTService):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -320,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
@@ -328,4 +328,4 @@ class FishAudioTTSService(InterruptibleTTSService):
|
|||||||
yield None
|
yield None
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -110,7 +110,7 @@ class GoogleImageGenService(ImageGenService):
|
|||||||
await self.stop_ttfb_metrics()
|
await self.stop_ttfb_metrics()
|
||||||
|
|
||||||
if not response or not response.generated_images:
|
if not response or not response.generated_images:
|
||||||
await self.push_error(error_msg="Image generation failed")
|
yield ErrorFrame("Image generation failed")
|
||||||
return
|
return
|
||||||
|
|
||||||
for img_response in response.generated_images:
|
for img_response in response.generated_images:
|
||||||
@@ -127,4 +127,4 @@ class GoogleImageGenService(ImageGenService):
|
|||||||
yield frame
|
yield frame
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(f"Image generation error: {str(e)}")
|
||||||
|
|||||||
@@ -737,7 +737,8 @@ class GoogleHttpTTSService(TTSService):
|
|||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
error_message = f"TTS generation error: {str(e)}"
|
||||||
|
yield ErrorFrame(error=error_message)
|
||||||
|
|
||||||
|
|
||||||
class GoogleBaseTTSService(TTSService):
|
class GoogleBaseTTSService(TTSService):
|
||||||
@@ -1244,4 +1245,5 @@ class GeminiTTSService(GoogleBaseTTSService):
|
|||||||
yield frame
|
yield frame
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
error_message = f"Gemini TTS generation error: {str(e)}"
|
||||||
|
yield ErrorFrame(error=error_message)
|
||||||
|
|||||||
@@ -146,6 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -161,7 +161,7 @@ class HeyGenClient:
|
|||||||
f"{self}::event_callback_task",
|
f"{self}::event_callback_task",
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Failed to setup HeyGenClient: {e}", exception=e)
|
logger.error(f"Failed to setup HeyGenClient: {e}")
|
||||||
await self.cleanup()
|
await self.cleanup()
|
||||||
|
|
||||||
async def cleanup(self) -> None:
|
async def cleanup(self) -> None:
|
||||||
@@ -179,7 +179,7 @@ class HeyGenClient:
|
|||||||
await self._task_manager.cancel_task(self._event_task)
|
await self._task_manager.cancel_task(self._event_task)
|
||||||
self._event_task = None
|
self._event_task = None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Exception during cleanup: {e}", exception=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.
|
||||||
@@ -229,7 +229,7 @@ class HeyGenClient:
|
|||||||
self._ws_receive_task_handler(), name="HeyGenClient_Websocket"
|
self._ws_receive_task_handler(), name="HeyGenClient_Websocket"
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"{self} initialization error: {e}", exception=e)
|
logger.error(f"{self} initialization error: {e}")
|
||||||
self._websocket = None
|
self._websocket = None
|
||||||
|
|
||||||
async def _ws_receive_task_handler(self):
|
async def _ws_receive_task_handler(self):
|
||||||
@@ -242,9 +242,7 @@ class HeyGenClient:
|
|||||||
except ConnectionClosedOK:
|
except ConnectionClosedOK:
|
||||||
break
|
break
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(
|
logger.error(f"Error processing WebSocket message: {e}")
|
||||||
error_msg=f"Error processing WebSocket message: {e}", exception=e
|
|
||||||
)
|
|
||||||
break
|
break
|
||||||
|
|
||||||
async def _handle_ws_server_event(self, event: dict) -> None:
|
async def _handle_ws_server_event(self, event: dict) -> None:
|
||||||
@@ -262,7 +260,7 @@ class HeyGenClient:
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"{self} disconnect error: {e}", exception=e)
|
logger.error(f"{self} disconnect error: {e}")
|
||||||
finally:
|
finally:
|
||||||
self._websocket = None
|
self._websocket = None
|
||||||
|
|
||||||
@@ -275,9 +273,7 @@ class HeyGenClient:
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
await self._websocket.send(json.dumps(message))
|
await self._websocket.send(json.dumps(message))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(
|
logger.error(f"Error sending message to HeyGen websocket: {e}")
|
||||||
error_msg=f"Error sending message to HeyGen websocket: {e}", exception=e
|
|
||||||
)
|
|
||||||
raise e
|
raise e
|
||||||
|
|
||||||
async def interrupt(self, event_id: str) -> None:
|
async def interrupt(self, event_id: str) -> None:
|
||||||
@@ -475,11 +471,9 @@ class HeyGenClient:
|
|||||||
await self._audio_frame_callback(audio_frame)
|
await self._audio_frame_callback(audio_frame)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(
|
logger.error(f"Error processing audio frame: {e}")
|
||||||
error_msg=f"Error processing audio frame: {e}", exception=e
|
|
||||||
)
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error processing audio frames: {e}", exception=e)
|
logger.error(f"Error processing audio frames: {e}")
|
||||||
finally:
|
finally:
|
||||||
logger.debug(f"Audio frame processing ended.")
|
logger.debug(f"Audio frame processing ended.")
|
||||||
|
|
||||||
@@ -506,11 +500,9 @@ class HeyGenClient:
|
|||||||
if self._transport_ready and self._video_frame_callback:
|
if self._transport_ready and self._video_frame_callback:
|
||||||
await self._video_frame_callback(image_frame)
|
await self._video_frame_callback(image_frame)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(
|
logger.error(f"Error processing individual video frame: {e}")
|
||||||
error_msg=f"Error processing individual video frame: {e}", exception=e
|
|
||||||
)
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error processing video frames: {e}", exception=e)
|
logger.error(f"Error processing video frames: {e}")
|
||||||
finally:
|
finally:
|
||||||
logger.debug(f"Video frame processing ended.")
|
logger.debug(f"Video frame processing ended.")
|
||||||
|
|
||||||
@@ -603,7 +595,7 @@ class HeyGenClient:
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"LiveKit initialization error: {e}", exception=e)
|
logger.error(f"LiveKit initialization error: {e}")
|
||||||
self._livekit_room = None
|
self._livekit_room = None
|
||||||
|
|
||||||
async def _livekit_disconnect(self):
|
async def _livekit_disconnect(self):
|
||||||
@@ -624,7 +616,7 @@ class HeyGenClient:
|
|||||||
self._livekit_room = None
|
self._livekit_room = None
|
||||||
logger.debug("Successfully disconnected from LiveKit room")
|
logger.debug("Successfully disconnected from LiveKit room")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"LiveKit disconnect error: {e}", exception=e)
|
logger.error(f"LiveKit disconnect error: {e}")
|
||||||
|
|
||||||
#
|
#
|
||||||
# Queue callback handling
|
# Queue callback handling
|
||||||
|
|||||||
@@ -299,11 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -264,7 +264,7 @@ class MiniMaxHttpTTSService(TTSService):
|
|||||||
) as response:
|
) as response:
|
||||||
if response.status != 200:
|
if response.status != 200:
|
||||||
error_message = f"MiniMax TTS error: HTTP {response.status}"
|
error_message = f"MiniMax TTS error: HTTP {response.status}"
|
||||||
await self.push_error(error_msg=error_message)
|
yield ErrorFrame(error=error_message)
|
||||||
return
|
return
|
||||||
|
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
|
|||||||
@@ -110,7 +110,7 @@ class MoondreamService(VisionService):
|
|||||||
if analysis fails.
|
if analysis fails.
|
||||||
"""
|
"""
|
||||||
if not self._model:
|
if not self._model:
|
||||||
await self.push_error(error_msg="Moondream model not initialized")
|
yield ErrorFrame("Moondream model not available")
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug(f"Analyzing image (bytes length: {len(frame.image)})")
|
logger.debug(f"Analyzing image (bytes length: {len(frame.image)})")
|
||||||
|
|||||||
@@ -363,14 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class NeuphonicHttpTTSService(TTSService):
|
class NeuphonicHttpTTSService(TTSService):
|
||||||
@@ -563,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
# Don't yield error frame for individual message failures
|
# Don't yield error frame for individual message failures
|
||||||
continue
|
continue
|
||||||
|
|
||||||
@@ -571,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
finally:
|
finally:
|
||||||
await self.stop_ttfb_metrics()
|
await self.stop_ttfb_metrics()
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -76,7 +76,7 @@ class OpenAIImageGenService(ImageGenService):
|
|||||||
image_url = image.data[0].url
|
image_url = image.data[0].url
|
||||||
|
|
||||||
if not image_url:
|
if not image_url:
|
||||||
await self.push_error(error_msg="Image generation failed")
|
yield ErrorFrame("Image generation failed")
|
||||||
return
|
return
|
||||||
|
|
||||||
# Load the image from the url
|
# Load the image from the url
|
||||||
|
|||||||
@@ -206,4 +206,4 @@ class OpenAITTSService(TTSService):
|
|||||||
yield frame
|
yield frame
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
except BadRequestError as e:
|
except BadRequestError as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -88,8 +88,8 @@ class PiperTTSService(TTSService):
|
|||||||
) as response:
|
) as response:
|
||||||
if response.status != 200:
|
if response.status != 200:
|
||||||
error = await response.text()
|
error = await response.text()
|
||||||
await self.push_error(
|
yield ErrorFrame(
|
||||||
error_msg=f"Error getting audio (status: {response.status}, error: {error})"
|
error=f"Error getting audio (status: {response.status}, error: {error})"
|
||||||
)
|
)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
|||||||
@@ -391,7 +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:
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
@@ -401,7 +401,7 @@ class PlayHTTTSService(InterruptibleTTSService):
|
|||||||
yield None
|
yield None
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class PlayHTHttpTTSService(TTSService):
|
class PlayHTHttpTTSService(TTSService):
|
||||||
@@ -621,7 +621,7 @@ class PlayHTHttpTTSService(TTSService):
|
|||||||
yield frame
|
yield frame
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
finally:
|
finally:
|
||||||
await self.stop_ttfb_metrics()
|
await self.stop_ttfb_metrics()
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -408,14 +408,14 @@ class RimeTTSService(AudioContextWordTTSService):
|
|||||||
await self._get_websocket().send(json.dumps(msg))
|
await self._get_websocket().send(json.dumps(msg))
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class RimeHttpTTSService(TTSService):
|
class RimeHttpTTSService(TTSService):
|
||||||
@@ -546,7 +546,7 @@ class RimeHttpTTSService(TTSService):
|
|||||||
) as response:
|
) as response:
|
||||||
if response.status != 200:
|
if response.status != 200:
|
||||||
error_message = f"Rime TTS error: HTTP {response.status}"
|
error_message = f"Rime TTS error: HTTP {response.status}"
|
||||||
await self.push_error(error_msg=error_message)
|
yield ErrorFrame(error=error_message)
|
||||||
return
|
return
|
||||||
|
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
@@ -563,7 +563,7 @@ class RimeHttpTTSService(TTSService):
|
|||||||
yield frame
|
yield frame
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
finally:
|
finally:
|
||||||
await self.stop_ttfb_metrics()
|
await self.stop_ttfb_metrics()
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -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:
|
||||||
await self.push_error(
|
yield ErrorFrame(f"Unexpected Riva response format: {str(ae)}")
|
||||||
error_msg=f"Unexpected response structure from Riva: {ae}", exception=ae
|
|
||||||
)
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating STT: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
|
|
||||||
class ParakeetSTTService(RivaSTTService):
|
class ParakeetSTTService(RivaSTTService):
|
||||||
|
|||||||
@@ -180,13 +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:
|
||||||
await self.push_error(error_msg=f"Timeout generating TTS: {text}")
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
|
||||||
return
|
|
||||||
except Exception as e:
|
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
|
||||||
yield TTSStoppedFrame()
|
|
||||||
return
|
|
||||||
|
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
|
|||||||
@@ -697,11 +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:
|
||||||
await self.push_error(error_msg=f"Error sending text: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
yield TTSStoppedFrame()
|
yield TTSStoppedFrame()
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
await self._connect()
|
await self._connect()
|
||||||
return
|
return
|
||||||
yield None
|
yield None
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating TTS: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -467,7 +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:
|
||||||
await self.push_error(error_msg=f"Error sending audio: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
await self._disconnect()
|
await self._disconnect()
|
||||||
|
|
||||||
def update_params(
|
def update_params(
|
||||||
|
|||||||
@@ -436,12 +436,12 @@ class UltravoxSTTService(AIService):
|
|||||||
yield LLMFullResponseEndFrame()
|
yield LLMFullResponseEndFrame()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error generating text: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
else:
|
else:
|
||||||
await self.push_error(error_msg="No model available for text generation")
|
yield ErrorFrame("No model available for text generation")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error processing audio: {e}", exception=e)
|
yield ErrorFrame(f"Error processing audio: {str(e)}")
|
||||||
finally:
|
finally:
|
||||||
self._buffer.is_processing = False
|
self._buffer.is_processing = False
|
||||||
self._buffer.frames = []
|
self._buffer.frames = []
|
||||||
|
|||||||
@@ -226,7 +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:
|
||||||
await self.push_error(exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|
||||||
async def _transcribe(self, audio: bytes) -> Transcription:
|
async def _transcribe(self, audio: bytes) -> Transcription:
|
||||||
"""Transcribe audio data to text.
|
"""Transcribe audio data to text.
|
||||||
|
|||||||
@@ -285,7 +285,7 @@ class WhisperSTTService(SegmentedSTTService):
|
|||||||
The service will normalize it to float32 in the range [-1, 1].
|
The service will normalize it to float32 in the range [-1, 1].
|
||||||
"""
|
"""
|
||||||
if not self._model:
|
if not self._model:
|
||||||
await self.push_error(error_msg="Whisper model not available")
|
yield ErrorFrame("Whisper model not available")
|
||||||
return
|
return
|
||||||
|
|
||||||
await self.start_processing_metrics()
|
await self.start_processing_metrics()
|
||||||
@@ -427,4 +427,4 @@ class WhisperSTTServiceMLX(WhisperSTTService):
|
|||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
await self.push_error(error_msg=f"Error processing audio: {e}", exception=e)
|
yield ErrorFrame(error=f"{self} error: {e}")
|
||||||
|
|||||||
@@ -181,9 +181,7 @@ 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()
|
||||||
await self.push_error(
|
yield ErrorFrame(error=f"Error getting audio (status: {r.status}, error: {text})")
|
||||||
error_msg=f"Error getting audio (status: {r.status}, error: {text})"
|
|
||||||
)
|
|
||||||
return
|
return
|
||||||
|
|
||||||
await self.start_tts_usage_metrics(text)
|
await self.start_tts_usage_metrics(text)
|
||||||
|
|||||||
Reference in New Issue
Block a user