@@ -50,6 +50,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
|
|
||||||
### Fixed
|
### Fixed
|
||||||
|
|
||||||
|
- Fixed an issue that would cause TTS websocket-based services to not cleanup
|
||||||
|
resources properly when disconnecting.
|
||||||
|
|
||||||
- Fixed a `TavusVideoService` issue that was causing audio choppiness.
|
- Fixed a `TavusVideoService` issue that was causing audio choppiness.
|
||||||
|
|
||||||
- Fixed an issue in `SmallWebRTCTransport` where an error was thrown if the
|
- Fixed an issue in `SmallWebRTCTransport` where an error was thrown if the
|
||||||
|
|||||||
@@ -197,7 +197,7 @@ class CartesiaTTSService(AudioContextWordTTSService):
|
|||||||
|
|
||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
logger.debug("Connecting to Cartesia")
|
logger.debug("Connecting to Cartesia")
|
||||||
self._websocket = await websockets.connect(
|
self._websocket = await websockets.connect(
|
||||||
@@ -215,11 +215,11 @@ class CartesiaTTSService(AudioContextWordTTSService):
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
logger.debug("Disconnecting from Cartesia")
|
logger.debug("Disconnecting from Cartesia")
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
|
|
||||||
self._context_id = None
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._context_id = None
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
def _get_websocket(self):
|
def _get_websocket(self):
|
||||||
if self._websocket:
|
if self._websocket:
|
||||||
@@ -279,7 +279,7 @@ class CartesiaTTSService(AudioContextWordTTSService):
|
|||||||
logger.debug(f"{self}: Generating TTS [{text}]")
|
logger.debug(f"{self}: Generating TTS [{text}]")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
|
|
||||||
if not self._context_id:
|
if not self._context_id:
|
||||||
|
|||||||
@@ -327,7 +327,7 @@ class ElevenLabsTTSService(InterruptibleWordTTSService):
|
|||||||
|
|
||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug("Connecting to ElevenLabs")
|
logger.debug("Connecting to ElevenLabs")
|
||||||
@@ -374,11 +374,11 @@ class ElevenLabsTTSService(InterruptibleWordTTSService):
|
|||||||
logger.debug("Disconnecting from ElevenLabs")
|
logger.debug("Disconnecting from ElevenLabs")
|
||||||
await self._websocket.send(json.dumps({"text": ""}))
|
await self._websocket.send(json.dumps({"text": ""}))
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
|
|
||||||
self._started = False
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._started = False
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
def _get_websocket(self):
|
def _get_websocket(self):
|
||||||
if self._websocket:
|
if self._websocket:
|
||||||
@@ -418,7 +418,7 @@ class ElevenLabsTTSService(InterruptibleWordTTSService):
|
|||||||
logger.debug(f"{self}: Generating TTS [{text}]")
|
logger.debug(f"{self}: Generating TTS [{text}]")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -116,7 +116,7 @@ class FishAudioTTSService(InterruptibleTTSService):
|
|||||||
|
|
||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug("Connecting to Fish Audio")
|
logger.debug("Connecting to Fish Audio")
|
||||||
@@ -141,16 +141,17 @@ class FishAudioTTSService(InterruptibleTTSService):
|
|||||||
stop_message = {"event": "stop"}
|
stop_message = {"event": "stop"}
|
||||||
await self._websocket.send(ormsgpack.packb(stop_message))
|
await self._websocket.send(ormsgpack.packb(stop_message))
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
self._request_id = None
|
|
||||||
self._started = False
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Error closing websocket: {e}")
|
logger.error(f"Error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._request_id = None
|
||||||
|
self._started = False
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
async def flush_audio(self):
|
async def flush_audio(self):
|
||||||
"""Flush any buffered audio by sending a flush event to Fish Audio."""
|
"""Flush any buffered audio by sending a flush event to Fish Audio."""
|
||||||
logger.trace(f"{self}: Flushing audio buffers")
|
logger.trace(f"{self}: Flushing audio buffers")
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
return
|
return
|
||||||
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))
|
||||||
|
|||||||
@@ -122,7 +122,7 @@ class LmntTTSService(InterruptibleTTSService):
|
|||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
"""Connect to LMNT websocket."""
|
"""Connect to LMNT websocket."""
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug("Connecting to LMNT")
|
logger.debug("Connecting to LMNT")
|
||||||
@@ -158,11 +158,11 @@ class LmntTTSService(InterruptibleTTSService):
|
|||||||
# errors on the websocket, so we just skip it for now.
|
# errors on the websocket, so we just skip it for now.
|
||||||
# await self._websocket.send(json.dumps({"eof": True}))
|
# await self._websocket.send(json.dumps({"eof": True}))
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
|
|
||||||
self._started = False
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._started = False
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
def _get_websocket(self):
|
def _get_websocket(self):
|
||||||
if self._websocket:
|
if self._websocket:
|
||||||
@@ -170,7 +170,7 @@ class LmntTTSService(InterruptibleTTSService):
|
|||||||
raise Exception("Websocket not connected")
|
raise Exception("Websocket not connected")
|
||||||
|
|
||||||
async def flush_audio(self):
|
async def flush_audio(self):
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
return
|
return
|
||||||
await self._get_websocket().send(json.dumps({"flush": True}))
|
await self._get_websocket().send(json.dumps({"flush": True}))
|
||||||
|
|
||||||
@@ -203,7 +203,7 @@ class LmntTTSService(InterruptibleTTSService):
|
|||||||
logger.debug(f"{self}: Generating TTS [{text}]")
|
logger.debug(f"{self}: Generating TTS [{text}]")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -175,6 +175,9 @@ class NeuphonicTTSService(InterruptibleTTSService):
|
|||||||
|
|
||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
try:
|
try:
|
||||||
|
if self._websocket and self._websocket.open:
|
||||||
|
return
|
||||||
|
|
||||||
logger.debug("Connecting to Neuphonic")
|
logger.debug("Connecting to Neuphonic")
|
||||||
|
|
||||||
tts_config = {
|
tts_config = {
|
||||||
@@ -203,11 +206,11 @@ class NeuphonicTTSService(InterruptibleTTSService):
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
logger.debug("Disconnecting from Neuphonic")
|
logger.debug("Disconnecting from Neuphonic")
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
|
|
||||||
self._started = False
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._started = False
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
async def _receive_messages(self):
|
async def _receive_messages(self):
|
||||||
async for message in self._websocket:
|
async for message in self._websocket:
|
||||||
@@ -235,7 +238,7 @@ class NeuphonicTTSService(InterruptibleTTSService):
|
|||||||
logger.debug(f"Generating TTS: [{text}]")
|
logger.debug(f"Generating TTS: [{text}]")
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -169,7 +169,7 @@ class PlayHTTTSService(InterruptibleTTSService):
|
|||||||
|
|
||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
|
|
||||||
logger.debug("Connecting to PlayHT")
|
logger.debug("Connecting to PlayHT")
|
||||||
@@ -197,11 +197,11 @@ class PlayHTTTSService(InterruptibleTTSService):
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
logger.debug("Disconnecting from PlayHT")
|
logger.debug("Disconnecting from PlayHT")
|
||||||
await self._websocket.close()
|
await self._websocket.close()
|
||||||
self._websocket = None
|
|
||||||
|
|
||||||
self._request_id = None
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._request_id = None
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
async def _get_websocket_url(self):
|
async def _get_websocket_url(self):
|
||||||
async with aiohttp.ClientSession() as session:
|
async with aiohttp.ClientSession() as session:
|
||||||
|
|||||||
@@ -182,7 +182,7 @@ class RimeTTSService(AudioContextWordTTSService):
|
|||||||
async def _connect_websocket(self):
|
async def _connect_websocket(self):
|
||||||
"""Connect to Rime websocket API with configured settings."""
|
"""Connect to Rime websocket API with configured settings."""
|
||||||
try:
|
try:
|
||||||
if self._websocket:
|
if self._websocket and self._websocket.open:
|
||||||
return
|
return
|
||||||
|
|
||||||
params = "&".join(f"{k}={v}" for k, v in self._settings.items())
|
params = "&".join(f"{k}={v}" for k, v in self._settings.items())
|
||||||
@@ -201,10 +201,11 @@ class RimeTTSService(AudioContextWordTTSService):
|
|||||||
if self._websocket:
|
if self._websocket:
|
||||||
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()
|
||||||
self._websocket = None
|
|
||||||
self._context_id = None
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error closing websocket: {e}")
|
logger.error(f"{self} error closing websocket: {e}")
|
||||||
|
finally:
|
||||||
|
self._context_id = None
|
||||||
|
self._websocket = None
|
||||||
|
|
||||||
def _get_websocket(self):
|
def _get_websocket(self):
|
||||||
"""Get active websocket connection or raise exception."""
|
"""Get active websocket connection or raise exception."""
|
||||||
@@ -316,7 +317,7 @@ class RimeTTSService(AudioContextWordTTSService):
|
|||||||
"""
|
"""
|
||||||
logger.debug(f"{self}: Generating TTS [{text}]")
|
logger.debug(f"{self}: Generating TTS [{text}]")
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
await self._connect()
|
await self._connect()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
|
|||||||
@@ -31,7 +31,7 @@ class WebsocketService(ABC):
|
|||||||
bool: True if connection is verified working, False otherwise
|
bool: True if connection is verified working, False otherwise
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
if not self._websocket:
|
if not self._websocket or self._websocket.closed:
|
||||||
return False
|
return False
|
||||||
await self._websocket.ping()
|
await self._websocket.ping()
|
||||||
return True
|
return True
|
||||||
|
|||||||
Reference in New Issue
Block a user