Merge pull request #3615 from omChauhanDev/fix/daily-transport-message-queue
fix(daily): queue outbound messages until transport joins
This commit is contained in:
1
changelog/3615.fixed.md
Normal file
1
changelog/3615.fixed.md
Normal file
@@ -0,0 +1 @@
|
|||||||
|
- Fixed race condition where `RTVIObserver` could send messages before `DailyTransport` join completed. Outbound messages are now queued & delivered after the transport is ready.
|
||||||
@@ -511,6 +511,7 @@ 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._join_message_queue: list = []
|
||||||
|
|
||||||
# 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
|
||||||
@@ -577,7 +578,8 @@ 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."
|
self._join_message_queue.append(frame)
|
||||||
|
return None
|
||||||
|
|
||||||
participant_id = None
|
participant_id = None
|
||||||
if isinstance(
|
if isinstance(
|
||||||
@@ -778,6 +780,8 @@ class DailyTransportClient(EventHandler):
|
|||||||
await self._callbacks.on_joined(data)
|
await self._callbacks.on_joined(data)
|
||||||
|
|
||||||
self._joined_event.set()
|
self._joined_event.set()
|
||||||
|
|
||||||
|
await self._flush_join_messages()
|
||||||
else:
|
else:
|
||||||
error_msg = f"Error joining {self._room_url}: {error}"
|
error_msg = f"Error joining {self._room_url}: {error}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
@@ -1551,6 +1555,12 @@ class DailyTransportClient(EventHandler):
|
|||||||
await callback(*args)
|
await callback(*args)
|
||||||
queue.task_done()
|
queue.task_done()
|
||||||
|
|
||||||
|
async def _flush_join_messages(self):
|
||||||
|
"""Send any messages that were queued before join completed."""
|
||||||
|
for frame in self._join_message_queue:
|
||||||
|
await self.send_message(frame)
|
||||||
|
self._join_message_queue.clear()
|
||||||
|
|
||||||
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