update keepalive times depending on watchdog timers
This commit is contained in:
@@ -428,9 +428,10 @@ class ElevenLabsTTSService(AudioContextWordTTSService):
|
|||||||
self._cumulative_time = word_times[-1][1]
|
self._cumulative_time = word_times[-1][1]
|
||||||
|
|
||||||
async def _keepalive_task_handler(self):
|
async def _keepalive_task_handler(self):
|
||||||
|
KEEPALIVE_SLEEP = 10 if self.watchdog_timers_enabled else 3
|
||||||
while True:
|
while True:
|
||||||
self.reset_watchdog()
|
self.reset_watchdog()
|
||||||
await asyncio.sleep(4)
|
await asyncio.sleep(KEEPALIVE_SLEEP)
|
||||||
try:
|
try:
|
||||||
if self._websocket and self._websocket.open:
|
if self._websocket and self._websocket.open:
|
||||||
if self._context_id:
|
if self._context_id:
|
||||||
|
|||||||
@@ -392,8 +392,8 @@ class GladiaSTTService(STTService):
|
|||||||
await self._send_buffered_audio()
|
await self._send_buffered_audio()
|
||||||
|
|
||||||
# Start tasks
|
# Start tasks
|
||||||
self._receive_task = asyncio.create_task(self._receive_task_handler())
|
self._receive_task = self.create_task(self._receive_task_handler())
|
||||||
self._keepalive_task = asyncio.create_task(self._keepalive_task_handler())
|
self._keepalive_task = self.create_task(self._keepalive_task_handler())
|
||||||
|
|
||||||
# Wait for tasks to complete
|
# Wait for tasks to complete
|
||||||
await asyncio.gather(self._receive_task, self._keepalive_task)
|
await asyncio.gather(self._receive_task, self._keepalive_task)
|
||||||
@@ -404,9 +404,9 @@ class GladiaSTTService(STTService):
|
|||||||
|
|
||||||
# Clean up tasks
|
# Clean up tasks
|
||||||
if self._receive_task:
|
if self._receive_task:
|
||||||
self._receive_task.cancel()
|
await self.cancel_task(self._receive_task)
|
||||||
if self._keepalive_task:
|
if self._keepalive_task:
|
||||||
self._keepalive_task.cancel()
|
await self.cancel_task(self._keepalive_task)
|
||||||
|
|
||||||
# Attempt reconnect using helper
|
# Attempt reconnect using helper
|
||||||
if not await self._maybe_reconnect():
|
if not await self._maybe_reconnect():
|
||||||
@@ -485,9 +485,11 @@ class GladiaSTTService(STTService):
|
|||||||
async def _keepalive_task_handler(self):
|
async def _keepalive_task_handler(self):
|
||||||
"""Send periodic empty audio chunks to keep the connection alive."""
|
"""Send periodic empty audio chunks to keep the connection alive."""
|
||||||
try:
|
try:
|
||||||
|
KEEPALIVE_SLEEP = 20 if self.watchdog_timers_enabled else 3
|
||||||
while self._connection_active:
|
while self._connection_active:
|
||||||
# Send keepalive every 20 seconds (Gladia times out after 30 seconds)
|
self.reset_watchdog()
|
||||||
await asyncio.sleep(20)
|
# Send keepalive (Gladia times out after 30 seconds)
|
||||||
|
await asyncio.sleep(KEEPALIVE_SLEEP)
|
||||||
if self._websocket and not self._websocket.closed:
|
if self._websocket and not self._websocket.closed:
|
||||||
# Send an empty audio chunk as keepalive
|
# Send an empty audio chunk as keepalive
|
||||||
empty_audio = b""
|
empty_audio = b""
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ from pipecat.processors.frame_processor import FrameDirection
|
|||||||
from pipecat.services.tts_service import InterruptibleTTSService, TTSService
|
from pipecat.services.tts_service import InterruptibleTTSService, TTSService
|
||||||
from pipecat.transcriptions.language import Language
|
from pipecat.transcriptions.language import Language
|
||||||
from pipecat.utils.tracing.service_decorators import traced_tts
|
from pipecat.utils.tracing.service_decorators import traced_tts
|
||||||
|
from pipecat.utils.watchdog_async_iterator import WatchdogAsyncIterator
|
||||||
|
|
||||||
try:
|
try:
|
||||||
import websockets
|
import websockets
|
||||||
@@ -221,7 +222,9 @@ class NeuphonicTTSService(InterruptibleTTSService):
|
|||||||
self._websocket = None
|
self._websocket = None
|
||||||
|
|
||||||
async def _receive_messages(self):
|
async def _receive_messages(self):
|
||||||
async for message in self._websocket:
|
async for message in WatchdogAsyncIterator(
|
||||||
|
self._websocket, reseter=self, watchdog_enabled=self.watchdog_timers_enabled
|
||||||
|
):
|
||||||
if isinstance(message, str):
|
if isinstance(message, str):
|
||||||
msg = json.loads(message)
|
msg = json.loads(message)
|
||||||
if msg.get("data", {}).get("audio") is not None:
|
if msg.get("data", {}).get("audio") is not None:
|
||||||
@@ -232,8 +235,10 @@ class NeuphonicTTSService(InterruptibleTTSService):
|
|||||||
await self.push_frame(frame)
|
await self.push_frame(frame)
|
||||||
|
|
||||||
async def _keepalive_task_handler(self):
|
async def _keepalive_task_handler(self):
|
||||||
|
KEEPALIVE_SLEEP = 10 if self.watchdog_timers_enabled else 3
|
||||||
while True:
|
while True:
|
||||||
await asyncio.sleep(10)
|
self.reset_watchdog()
|
||||||
|
await asyncio.sleep(KEEPALIVE_SLEEP)
|
||||||
await self._send_text("")
|
await self._send_text("")
|
||||||
|
|
||||||
async def _send_text(self, text: str):
|
async def _send_text(self, text: str):
|
||||||
|
|||||||
Reference in New Issue
Block a user