Merge pull request #585 from pipecat-ai/aleix/daily-urgent-transport-message-hang
transports(daily): send transport messages in a task
This commit is contained in:
@@ -13,6 +13,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
|
|||||||
`filter_code` to filter code from text and `filter_tables` to filter tables
|
`filter_code` to filter code from text and `filter_tables` to filter tables
|
||||||
from text.
|
from text.
|
||||||
|
|
||||||
|
### Fixed
|
||||||
|
|
||||||
|
- Fixed an issue in Daily transport that would cause tasks to be hanging if
|
||||||
|
urgent transport messages were being sent from a transport event handler.
|
||||||
|
|
||||||
## [0.0.43] - 2024-10-10
|
## [0.0.43] - 2024-10-10
|
||||||
|
|
||||||
### Added
|
### Added
|
||||||
|
|||||||
@@ -720,15 +720,27 @@ class DailyOutputTransport(BaseOutputTransport):
|
|||||||
|
|
||||||
self._client = client
|
self._client = client
|
||||||
|
|
||||||
|
# Task to process outgoing messages.
|
||||||
|
self._messages_task = None
|
||||||
|
self._messages_queue = asyncio.Queue()
|
||||||
|
|
||||||
async def start(self, frame: StartFrame):
|
async def start(self, frame: StartFrame):
|
||||||
# Parent start.
|
# Parent start.
|
||||||
await super().start(frame)
|
await super().start(frame)
|
||||||
# Join the room.
|
# Join the room.
|
||||||
await self._client.join()
|
await self._client.join()
|
||||||
|
# Start messages task
|
||||||
|
self._messages_task = self.get_event_loop().create_task(self._messages_task_handler())
|
||||||
|
|
||||||
async def stop(self, frame: EndFrame):
|
async def stop(self, frame: EndFrame):
|
||||||
# Parent stop.
|
# Parent stop.
|
||||||
await super().stop(frame)
|
await super().stop(frame)
|
||||||
|
# Cancel messages task
|
||||||
|
if self._messages_task:
|
||||||
|
self._messages_task.cancel()
|
||||||
|
await self._messages_task
|
||||||
|
self._messages_task = None
|
||||||
|
self._messages_task = None
|
||||||
# Leave the room.
|
# Leave the room.
|
||||||
await self._client.leave()
|
await self._client.leave()
|
||||||
|
|
||||||
@@ -743,7 +755,7 @@ class DailyOutputTransport(BaseOutputTransport):
|
|||||||
await self._client.cleanup()
|
await self._client.cleanup()
|
||||||
|
|
||||||
async def send_message(self, frame: TransportMessageFrame | TransportMessageUrgentFrame):
|
async def send_message(self, frame: TransportMessageFrame | TransportMessageUrgentFrame):
|
||||||
await self._client.send_message(frame)
|
await self._messages_queue.put(frame)
|
||||||
|
|
||||||
async def send_metrics(self, frame: MetricsFrame):
|
async def send_metrics(self, frame: MetricsFrame):
|
||||||
metrics = {}
|
metrics = {}
|
||||||
@@ -768,7 +780,7 @@ class DailyOutputTransport(BaseOutputTransport):
|
|||||||
message = DailyTransportMessageFrame(
|
message = DailyTransportMessageFrame(
|
||||||
message={"type": "pipecat-metrics", "metrics": metrics}
|
message={"type": "pipecat-metrics", "metrics": metrics}
|
||||||
)
|
)
|
||||||
await self._client.send_message(message)
|
await self._messages_queue.put(message)
|
||||||
|
|
||||||
async def write_raw_audio_frames(self, frames: bytes):
|
async def write_raw_audio_frames(self, frames: bytes):
|
||||||
await self._client.write_raw_audio_frames(frames)
|
await self._client.write_raw_audio_frames(frames)
|
||||||
@@ -776,6 +788,17 @@ class DailyOutputTransport(BaseOutputTransport):
|
|||||||
async def write_frame_to_camera(self, frame: OutputImageRawFrame):
|
async def write_frame_to_camera(self, frame: OutputImageRawFrame):
|
||||||
await self._client.write_frame_to_camera(frame)
|
await self._client.write_frame_to_camera(frame)
|
||||||
|
|
||||||
|
async def _messages_task_handler(self):
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
message = await self._messages_queue.get()
|
||||||
|
await self._client.send_message(message)
|
||||||
|
self._messages_queue.task_done()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
break
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception(f"{self} error processing message queue: {e}")
|
||||||
|
|
||||||
|
|
||||||
class DailyTransport(BaseTransport):
|
class DailyTransport(BaseTransport):
|
||||||
def __init__(
|
def __init__(
|
||||||
|
|||||||
Reference in New Issue
Block a user