fix(daily): queue outbound messages until transport joins
This commit is contained in:
@@ -501,6 +501,8 @@ class DailyTransportClient(EventHandler):
|
|||||||
self._event_task = None
|
self._event_task = None
|
||||||
self._audio_task = None
|
self._audio_task = None
|
||||||
self._video_task = None
|
self._video_task = None
|
||||||
|
self._message_queue: Optional[asyncio.Queue] = None
|
||||||
|
self._message_task = None
|
||||||
|
|
||||||
# Input and ouput sample rates. They will be initialize on setup().
|
# Input and ouput sample rates. They will be initialize on setup().
|
||||||
self._in_sample_rate = 0
|
self._in_sample_rate = 0
|
||||||
@@ -567,7 +569,9 @@ class DailyTransportClient(EventHandler):
|
|||||||
error: An error description or None.
|
error: An error description or None.
|
||||||
"""
|
"""
|
||||||
if not self._joined:
|
if not self._joined:
|
||||||
return "Unable to send messages before joining."
|
if self._message_queue:
|
||||||
|
await self._message_queue.put((frame,))
|
||||||
|
return None
|
||||||
|
|
||||||
participant_id = None
|
participant_id = None
|
||||||
if isinstance(
|
if isinstance(
|
||||||
@@ -673,6 +677,12 @@ class DailyTransportClient(EventHandler):
|
|||||||
f"{self}::event_callback_task",
|
f"{self}::event_callback_task",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
self._message_queue = asyncio.Queue()
|
||||||
|
self._message_task = self._task_manager.create_task(
|
||||||
|
self._message_task_handler(self._message_queue),
|
||||||
|
f"{self}::message_task",
|
||||||
|
)
|
||||||
|
|
||||||
async def cleanup(self):
|
async def cleanup(self):
|
||||||
"""Cleanup client resources and cancel tasks."""
|
"""Cleanup client resources and cancel tasks."""
|
||||||
if self._event_task and self._task_manager:
|
if self._event_task and self._task_manager:
|
||||||
@@ -684,6 +694,9 @@ class DailyTransportClient(EventHandler):
|
|||||||
if self._video_task and self._task_manager:
|
if self._video_task and self._task_manager:
|
||||||
await self._task_manager.cancel_task(self._video_task)
|
await self._task_manager.cancel_task(self._video_task)
|
||||||
self._video_task = None
|
self._video_task = None
|
||||||
|
if self._message_task and self._task_manager:
|
||||||
|
await self._task_manager.cancel_task(self._message_task)
|
||||||
|
self._message_task = None
|
||||||
# Make sure we don't block the event loop in case `client.release()`
|
# Make sure we don't block the event loop in case `client.release()`
|
||||||
# takes extra time.
|
# takes extra time.
|
||||||
await self._get_event_loop().run_in_executor(self._executor, self._cleanup)
|
await self._get_event_loop().run_in_executor(self._executor, self._cleanup)
|
||||||
@@ -1541,6 +1554,14 @@ class DailyTransportClient(EventHandler):
|
|||||||
await callback(*args)
|
await callback(*args)
|
||||||
queue.task_done()
|
queue.task_done()
|
||||||
|
|
||||||
|
async def _message_task_handler(self, queue: asyncio.Queue):
|
||||||
|
"""Handle queued messages after transport is joined."""
|
||||||
|
while True:
|
||||||
|
await self._joined_event.wait()
|
||||||
|
(frame,) = await queue.get()
|
||||||
|
await self.send_message(frame)
|
||||||
|
queue.task_done()
|
||||||
|
|
||||||
def _get_event_loop(self) -> asyncio.AbstractEventLoop:
|
def _get_event_loop(self) -> asyncio.AbstractEventLoop:
|
||||||
"""Get the event loop from the task manager."""
|
"""Get the event loop from the task manager."""
|
||||||
if not self._task_manager:
|
if not self._task_manager:
|
||||||
|
|||||||
Reference in New Issue
Block a user